adds basic web server protections (static-threaded); fixes #3291; also some improvements on the exit strategy of netdata
Costa Tsaousis (ktsaou) committed
Jan 15, 2018 at 23:57 UTC
98d08bc77c815a01bbda6f29af64423062872aff
18 files changed
+279
-179
src/backends.c
+4
-4
@@ -500,11 +500,11 @@ inline uint32_t backend_parse_data_source(const char *source, uint32_t mode) {
500
501
static void backends_main_cleanup(void *ptr) {
502
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
503
- if(static_thread->enabled) {
504
- info("cleaning up...");
503
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
504
506
- static_thread->enabled = 0;
507
- }
505
+ info("cleaning up...");
506
+
507
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
508
}
509
510
void *backends_main(void *ptr) {
src/health.c
+4
-4
@@ -340,11 +340,11 @@ static inline int check_if_resumed_from_suspention(void) {
340
341
static void health_main_cleanup(void *ptr) {
342
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
343
- if(static_thread->enabled) {
344
- info("cleaning up...");
343
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
344
346
- static_thread->enabled = 0;
347
- }
345
+ info("cleaning up...");
346
+
347
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
348
}
349
350
void *health_main(void *ptr) {
src/main.c
+10
-15
@@ -3,29 +3,24 @@
3
extern void *cgroups_main(void *ptr);
4
5
void netdata_cleanup_and_exit(int ret) {
6
- netdata_exit = 1;
6
+ // enabling this, is wrong
7
+ // because the threads will be cancelled while cleaning up
8
+ // netdata_exit = 1;
9
10
error_log_limit_unlimited();
11
info("EXIT: netdata prepares to exit with code %d...", ret);
12
11
- if(ret) {
12
- // this is bad - exiting due to a fatal condition
13
+ // cleanup/save the database and exit
14
+ info("EXIT: cleaning up the database...");
15
+ rrdhost_cleanup_all();
16
14
- // cleanup/save the database and exit
15
- info("EXIT: cleaning up the database...");
16
- rrdhost_cleanup_all();
17
- }
18
- else {
17
+ if(!ret) {
18
// exit cleanly
19
20
// stop everything
21
info("EXIT: stopping master threads...");
22
cancel_main_threads();
23
25
- // cleanup the database (delete files not needed)
26
- info("EXIT: cleaning up the database...");
27
- rrdhost_cleanup_all();
28
-
24
// free the database
25
info("EXIT: freeing database memory...");
26
rrdhost_free_all();
@@ -198,7 +193,7 @@ void cancel_main_threads() {
193
194
int i, found = 0, max = 5 * USEC_PER_SEC, step = 100000;
195
for (i = 0; static_threads[i].name != NULL ; i++) {
201
- if(static_threads[i].enabled) {
196
+ if(static_threads[i].enabled == NETDATA_MAIN_THREAD_RUNNING) {
197
info("EXIT: Stopping master thread: %s", static_threads[i].name);
198
netdata_thread_cancel(*static_threads[i].thread);
199
found++;
@@ -211,14 +206,14 @@ void cancel_main_threads() {
206
sleep_usec(step);
207
found = 0;
208
for (i = 0; static_threads[i].name != NULL ; i++) {
214
- if (static_threads[i].enabled)
209
+ if (static_threads[i].enabled != NETDATA_MAIN_THREAD_EXITED)
210
found++;
211
}
212
}
213
214
if(found) {
215
for (i = 0; static_threads[i].name != NULL ; i++) {
221
- if (static_threads[i].enabled)
216
+ if (static_threads[i].enabled != NETDATA_MAIN_THREAD_EXITED)
217
error("Master thread %s takes too long to exit. Giving up...", static_threads[i].name);
218
}
219
}
src/main.h
+4
@@ -1,6 +1,10 @@
1
#ifndef NETDATA_MAIN_H
2
#define NETDATA_MAIN_H 1
3
4
+#define NETDATA_MAIN_THREAD_RUNNING CONFIG_BOOLEAN_YES
5
+#define NETDATA_MAIN_THREAD_EXITING (CONFIG_BOOLEAN_YES + 1)
6
+#define NETDATA_MAIN_THREAD_EXITED CONFIG_BOOLEAN_NO
7
+
8
/**
9
* This struct contains information about command line options.
10
*/
src/plugin_checks.c
+4
-4
@@ -4,11 +4,11 @@
4
5
static void checks_main_cleanup(void *ptr) {
6
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
7
- if(static_thread->enabled) {
8
- info("cleaning up...");
7
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
8
10
- static_thread->enabled = 0;
11
- }
9
+ info("cleaning up...");
10
+
11
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
12
}
13
14
void *checks_main(void *ptr) {
src/plugin_freebsd.c
+5
-5
@@ -68,11 +68,11 @@ static struct freebsd_module {
68
69
static void freebsd_main_cleanup(void *ptr) {
70
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
71
- if(static_thread->enabled) {
72
- info("cleaning up...");
71
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
72
74
- static_thread->enabled = 0;
75
- }
73
+ info("cleaning up...");
74
+
75
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
76
}
77
78
void *freebsd_main(void *ptr) {
@@ -82,7 +82,7 @@ void *freebsd_main(void *ptr) {
82
83
// initialize FreeBSD plugin
84
if (freebsd_plugin_init())
85
- netdata_exit = 1;
85
+ netdata_cleanup_and_exit(1);
86
87
// check the enabled status for each module
88
int i;
src/plugin_idlejitter.c
+4
-4
@@ -4,11 +4,11 @@
4
5
static void cpuidlejitter_main_cleanup(void *ptr) {
6
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
7
- if(static_thread->enabled) {
8
- info("cleaning up...");
7
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
8
10
- static_thread->enabled = 0;
11
- }
9
+ info("cleaning up...");
10
+
11
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
12
}
13
14
void *cpuidlejitter_main(void *ptr) {
src/plugin_macos.c
+4
-4
@@ -2,11 +2,11 @@
2
3
static void macos_main_cleanup(void *ptr) {
4
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
5
- if(static_thread->enabled) {
6
- info("cleaning up...");
5
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
6
8
- static_thread->enabled = 0;
9
- }
7
+ info("cleaning up...");
8
+
9
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
10
}
11
12
void *macos_main(void *ptr) {
src/plugin_nfacct.c
+5
-6
@@ -753,19 +753,18 @@ static void nfacct_send_metrics() {
753
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
- info("cleaning up...");
756
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
757
+ info("cleaning up...");
758
759
#ifdef DO_NFACCT
760
- nfacct_cleanup();
760
+ nfacct_cleanup();
761
#endif
762
763
#ifdef DO_NFSTAT
764
- nfstat_cleanup();
764
+ nfstat_cleanup();
765
#endif
766
767
- static_thread->enabled = 0;
768
- }
767
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
768
}
769
770
void *nfacct_main(void *ptr) {
src/plugin_proc.c
+4
-4
@@ -66,11 +66,11 @@ static struct proc_module {
66
67
static void proc_main_cleanup(void *ptr) {
68
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
69
- if(static_thread->enabled) {
70
- info("cleaning up...");
69
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
70
72
- static_thread->enabled = 0;
73
- }
71
+ info("cleaning up...");
72
+
73
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
74
}
75
76
void *proc_main(void *ptr) {
src/plugin_proc_diskspace.c
+4
-4
@@ -330,11 +330,11 @@ static inline void do_disk_space_stats(struct mountinfo *mi, int update_every) {
330
331
static void diskspace_main_cleanup(void *ptr) {
332
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
333
- if(static_thread->enabled) {
334
- info("cleaning up...");
333
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
334
336
- static_thread->enabled = 0;
337
- }
335
+ info("cleaning up...");
336
+
337
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
338
}
339
340
void *proc_diskspace_main(void *ptr) {
src/plugin_tc.c
+12
-12
@@ -832,24 +832,24 @@ static pid_t tc_child_pid = 0;
832
833
static void tc_main_cleanup(void *ptr) {
834
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
835
- if(static_thread->enabled) {
836
- info("cleaning up...");
835
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
836
838
- if(tc_child_pid) {
839
- info("TC: killing with SIGTERM tc-qos-helper process %d", tc_child_pid);
840
- if(killpid(tc_child_pid, SIGTERM) != -1) {
841
- siginfo_t info;
837
+ info("cleaning up...");
838
843
- info("TC: waiting for tc plugin child process pid %d to exit...", tc_child_pid);
844
- waitid(P_PID, (id_t) tc_child_pid, &info, WEXITED);
845
- // info("TC: finished tc plugin child process pid %d.", tc_child_pid);
846
- }
839
+ if(tc_child_pid) {
840
+ info("TC: killing with SIGTERM tc-qos-helper process %d", tc_child_pid);
841
+ if(killpid(tc_child_pid, SIGTERM) != -1) {
842
+ siginfo_t info;
843
848
- tc_child_pid = 0;
844
+ info("TC: waiting for tc plugin child process pid %d to exit...", tc_child_pid);
845
+ waitid(P_PID, (id_t) tc_child_pid, &info, WEXITED);
846
+ // info("TC: finished tc plugin child process pid %d.", tc_child_pid);
847
}
848
851
- static_thread->enabled = 0;
849
+ tc_child_pid = 0;
850
}
851
+
852
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
853
}
854
855
void *tc_main(void *ptr) {
src/plugins_d.c
+11
-12
@@ -568,20 +568,19 @@ void *pluginsd_worker_thread(void *arg) {
568
569
static void pluginsd_main_cleanup(void *data) {
570
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)data;
571
- if(static_thread->enabled) {
572
- info("cleaning up...");
573
-
574
- struct plugind *cd;
575
- for (cd = pluginsd_root; cd; cd = cd->next) {
576
- if (cd->enabled && !cd->obsolete) {
577
- info("stopping plugin thread: %s", cd->id);
578
- netdata_thread_cancel(cd->thread);
579
- }
571
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
572
+ info("cleaning up...");
573
+
574
+ struct plugind *cd;
575
+ for (cd = pluginsd_root; cd; cd = cd->next) {
576
+ if (cd->enabled && !cd->obsolete) {
577
+ info("stopping plugin thread: %s", cd->id);
578
+ netdata_thread_cancel(cd->thread);
579
}
581
-
582
- info("cleanup completed.");
583
- static_thread->enabled = 0;
580
}
581
+
582
+ info("cleanup completed.");
583
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
584
}
585
586
void *pluginsd_main(void *ptr) {
src/socket.c
+83
-14
@@ -1035,6 +1035,14 @@ inline POLLINFO *poll_add_fd(POLLJOB *p
1035
pi->rcv_callback = rcv_callback;
1036
pi->snd_callback = snd_callback;
1037
1038
+ pi->connected_t = now_boottime_sec();
1039
+ pi->last_received_t = 0;
1040
+ pi->last_sent_t = 0;
1041
+ pi->last_sent_t = 0;
1042
+ pi->recv_count = 0;
1043
+ pi->send_count = 0;
1044
+
1045
+ netdata_thread_disable_cancelability();
1046
p->used++;
1047
if(unlikely(pi->slot > p->max))
1048
p->max = pi->slot;
@@ -1046,6 +1054,7 @@ inline POLLINFO *poll_add_fd(POLLJOB *p
1054
if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
1055
p->min = pi->slot;
1056
}
1057
+ netdata_thread_enable_cancelability();
1058
1059
debug(D_POLLFD, "POLLFD: ADD: completed, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
1060
@@ -1060,13 +1069,15 @@ inline void poll_close_fd(POLLINFO *pi) {
1069
1070
if(unlikely(pf->fd == -1)) return;
1071
1072
+ netdata_thread_disable_cancelability();
1073
+
1074
if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1075
pi->del_callback(pi);
1065
- }
1076
1067
- if(likely(!(pi->flags & POLLINFO_FLAG_DONT_CLOSE))) {
1068
- if(close(pf->fd) == -1)
1069
- error("Failed to close() poll_events() socket %d", pf->fd);
1077
+ if(likely(!(pi->flags & POLLINFO_FLAG_DONT_CLOSE))) {
1078
+ if(close(pf->fd) == -1)
1079
+ error("Failed to close() poll_events() socket %d", pf->fd);
1080
+ }
1081
}
1082
1083
pf->fd = -1;
@@ -1102,6 +1113,7 @@ inline void poll_close_fd(POLLINFO *pi) {
1113
}
1114
}
1115
}
1116
+ netdata_thread_enable_cancelability();
1117
1118
debug(D_POLLFD, "POLLFD: DEL: completed, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
1119
}
@@ -1164,7 +1176,7 @@ static void poll_events_cleanup(void *data) {
1176
freez(p->inf);
1177
}
1178
1167
-static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, short int revents) {
1179
+static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, short int revents, time_t now) {
1180
short int events = pf->events;
1181
int fd = pf->fd;
1182
pf->revents = 0;
@@ -1180,6 +1192,9 @@ static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, sho
1192
if(revents & POLLIN || revents & POLLPRI) {
1193
// receiving data
1194
1195
+ pi->last_received_t = now;
1196
+ pi->recv_count++;
1197
+
1198
if(likely(pi->flags & POLLINFO_FLAG_SERVER_SOCKET)) {
1199
// new connection
1200
// debug(D_POLLFD, "POLLFD: LISTENER: accepting connections from slot %zu (fd %d)", i, fd);
@@ -1259,6 +1274,15 @@ static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, sho
1274
poll_close_fd(pi);
1275
return;
1276
}
1277
+
1278
+#ifdef NETDATA_INTERNAL_CHECKS
1279
+ // this is common - it is used for web server file copies
1280
+ if(!(pf->events & (POLLIN|POLLOUT))) {
1281
+ error("POLLFD: LISTENER: after reading, client slot %zu (fd %d) from '%s:%s' was left without expecting input or output. ", i, fd, pi->client_ip?pi->client_ip:"<undefined-ip>", pi->client_port?pi->client_port:"<undefined-port>");
1282
+ //poll_close_fd(pi);
1283
+ //return;
1284
+ }
1285
+#endif
1286
}
1287
}
1288
@@ -1266,11 +1290,20 @@ static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, sho
1290
// sending data
1291
debug(D_POLLFD, "POLLFD: LISTENER: sending data to socket on slot %zu (fd %d)", i, fd);
1292
1293
+ pi->last_sent_t = now;
1294
+ pi->send_count++;
1295
+
1296
pf->events = 0;
1297
if (pi->snd_callback(pi, &pf->events) == -1) {
1298
poll_close_fd(pi);
1299
return;
1300
}
1301
+
1302
+ if(unlikely(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET && !(pf->events & (POLLIN|POLLOUT)))) {
1303
+ error("POLLFD: LISTENER: after sending, client slot %zu (fd %d) from '%s:%s' was left without expecting input or output - closing it. ", i, fd, pi->client_ip?pi->client_ip:"<undefined-ip>", pi->client_port?pi->client_port:"<undefined-port>");
1304
+ poll_close_fd(pi);
1305
+ return;
1306
+ }
1307
}
1308
1309
if(unlikely(revents & POLLERR)) {
@@ -1318,6 +1351,10 @@ void poll_events(LISTEN_SOCKETS *sockets
1351
.inf = NULL,
1352
.first_free = NULL,
1353
1354
+ .complete_request_timeout = 30,
1355
+ .idle_timeout = web_client_timeout,
1356
+ .checks_every = (web_client_timeout / 3) + 1,
1357
+
1358
.access_list = access_list,
1359
1360
.add_callback = add_callback?add_callback:poll_default_add_callback,
@@ -1346,13 +1383,15 @@ void poll_events(LISTEN_SOCKETS *sockets
1383
info("POLLFD: LISTENER: listening on '%s'", (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN");
1384
}
1385
1349
- int timeout = -1; // wait forever
1386
+ int timeout = 1; // wait forever
1387
+ time_t last_check = now_boottime_sec();
1388
1389
netdata_thread_cleanup_push(poll_events_cleanup, &p);
1390
1391
while(!netdata_exit) {
1392
debug(D_POLLFD, "POLLFD: LISTENER: Waiting on %zu sockets...", p.max + 1);
1393
retval = poll(p.fds, p.max + 1, timeout);
1394
+ time_t now = now_boottime_sec();
1395
1396
if(unlikely(retval == -1)) {
1397
error("POLLFD: LISTENER: poll() failed while waiting on %zu sockets.", p.max + 1);
@@ -1360,16 +1399,46 @@ void poll_events(LISTEN_SOCKETS *sockets
1399
}
1400
else if(unlikely(!retval)) {
1401
debug(D_POLLFD, "POLLFD: LISTENER: poll() timeout.");
1363
- continue;
1402
+ }
1403
+ else {
1404
+ for (i = 0; i <= p.max; i++) {
1405
+ struct pollfd *pf = &p.fds[i];
1406
+ short int revents = pf->revents;
1407
+ if (unlikely(revents))
1408
+ poll_events_process(&p, &p.inf[i], pf, revents, now);
1409
+ }
1410
}
1411
1366
- if(unlikely(netdata_exit)) break;
1367
-
1368
- for(i = 0 ; i <= p.max ; i++) {
1369
- struct pollfd *pf = &p.fds[i];
1370
- short int revents = pf->revents;
1371
- if(unlikely(revents))
1372
- poll_events_process(&p, &p.inf[i], pf, revents);
1412
+ if(now - last_check > p.checks_every) {
1413
+ // security checks
1414
+ for(i = 0; i <= p.max; i++) {
1415
+ POLLINFO *pi = &p.inf[i];
1416
+
1417
+ if(likely(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET)) {
1418
+ if (unlikely(pi->send_count == 0 && (now - pi->connected_t) >= p.complete_request_timeout)) {
1419
+ errno = 0;
1420
+ error("POLLFD: LISTENER: client slot %zu (fd %d) from '%s:%s' has not sent a complete request in %zu seconds - closing it. "
1421
+ , i
1422
+ , pi->fd
1423
+ , pi->client_ip ? pi->client_ip : "<undefined-ip>"
1424
+ , pi->client_port ? pi->client_port : "<undefined-port>"
1425
+ , (size_t) p.complete_request_timeout
1426
+ );
1427
+ poll_close_fd(pi);
1428
+ }
1429
+ else if(unlikely(pi->recv_count && now - ((pi->last_received_t > pi->last_sent_t) ? pi->last_received_t : pi->last_sent_t) >= p.idle_timeout )) {
1430
+ errno = 0;
1431
+ error("POLLFD: LISTENER: client slot %zu (fd %d) from '%s:%s' is idle for more than %zu seconds - closing it. "
1432
+ , i
1433
+ , pi->fd
1434
+ , pi->client_ip ? pi->client_ip : "<undefined-ip>"
1435
+ , pi->client_port ? pi->client_port : "<undefined-port>"
1436
+ , (size_t) p.idle_timeout
1437
+ );
1438
+ poll_close_fd(pi);
1439
+ }
1440
+ }
1441
+ }
1442
}
1443
}
1444
src/socket.h
+19
-7
@@ -63,15 +63,22 @@ extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flag
63
typedef struct poll POLLJOB;
64
65
typedef struct pollinfo {
66
- POLLJOB *p; // the parent
67
- size_t slot; // the slot id
66
+ POLLJOB *p; // the parent
67
+ size_t slot; // the slot id
68
69
- int fd; // the file descriptor
70
- int socktype; // the client socket type
71
- char *client_ip; // the connected client IP
72
- char *client_port; // the connected client port
69
+ int fd; // the file descriptor
70
+ int socktype; // the client socket type
71
+ char *client_ip; // the connected client IP
72
+ char *client_port; // the connected client port
73
74
- uint32_t flags; // internal flags
74
+ time_t connected_t; // the time the socket connected
75
+ time_t last_received_t; // the time the socket last received data
76
+ time_t last_sent_t; // the time the socket last sent data
77
+
78
+ size_t recv_count; // the number of times the socket was ready for inbound traffic
79
+ size_t send_count; // the number of times the socket was ready for outbound traffic
80
+
81
+ uint32_t flags; // internal flags
82
83
// callbacks for this socket
84
void (*del_callback)(struct pollinfo *pi);
@@ -93,6 +100,11 @@ struct poll {
100
size_t used;
101
size_t min;
102
size_t max;
103
+
104
+ time_t complete_request_timeout;
105
+ time_t idle_timeout;
106
+ time_t checks_every;
107
+
108
struct pollfd *fds;
109
struct pollinfo *inf;
110
struct pollinfo *first_free;
src/statsd.c
+36
-25
@@ -254,6 +254,7 @@ static struct statsd {
254
char *histogram_percentile_str;
255
256
int threads;
257
+ int *collection_threads_status;
258
netdata_thread_t *collection_threads;
259
260
LISTEN_SOCKETS sockets;
@@ -323,6 +324,7 @@ static struct statsd {
324
.histogram_percentile = 95.0,
325
.histogram_increase_step = 10,
326
.threads = 0,
327
+ .collection_threads_status = NULL,
328
.collection_threads = NULL,
329
.sockets = {
330
.config_section = CONFIG_SECTION_STATSD,
@@ -684,6 +686,7 @@ struct statsd_tcp {
686
687
#ifdef HAVE_RECVMMSG
688
struct statsd_udp {
689
+ int *running;
690
STATSD_SOCKET_DATA_TYPE type;
691
size_t size;
692
struct iovec *iovecs;
@@ -691,6 +694,7 @@ struct statsd_udp {
694
};
695
#else
696
struct statsd_udp {
697
+ int *running;
698
STATSD_SOCKET_DATA_TYPE type;
699
char buffer[STATSD_UDP_BUFFER_SIZE];
700
};
@@ -875,31 +879,32 @@ static int statsd_snd_callback(POLLINFO *pi, short int *events) {
879
// statsd child thread to collect metrics from network
880
881
void statsd_collector_thread_cleanup(void *data) {
878
- static __thread int executed = 0;
879
- if(!executed) {
880
- executed = 1;
881
- info("cleaning up...");
882
+ struct statsd_udp *d = data;
883
+ *d->running = 0;
884
883
- struct statsd_udp *d = data;
885
+ info("cleaning up...");
886
887
#ifdef HAVE_RECVMMSG
886
- size_t i;
887
- for (i = 0; i < d->size; i++)
888
- freez(d->iovecs[i].iov_base);
888
+ size_t i;
889
+ for (i = 0; i < d->size; i++)
890
+ freez(d->iovecs[i].iov_base);
891
890
- freez(d->iovecs);
891
- freez(d->msgs);
892
+ freez(d->iovecs);
893
+ freez(d->msgs);
894
#endif
895
894
- freez(d);
895
- }
896
+ freez(d);
897
}
898
899
void *statsd_collector_thread(void *ptr) {
899
- int id = *((int *)ptr);
900
- info("STATSD collector thread No %d created with task id %d", id + 1, gettid());
900
+ int *running = (int *)ptr;
901
+ *running = 1;
902
+
903
+ info("STATSD collector thread started with taskid %d", gettid());
904
905
struct statsd_udp *d = callocz(sizeof(struct statsd_udp), 1);
906
+ d->running = running;
907
+
908
netdata_thread_cleanup_push(statsd_collector_thread_cleanup, d);
909
910
#ifdef HAVE_RECVMMSG
@@ -1969,23 +1974,27 @@ static int statsd_listen_sockets_setup(void) {
1974
1975
static void statsd_main_cleanup(void *data) {
1976
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)data;
1972
- if(static_thread->enabled) {
1973
- info("cleaning up...");
1977
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
1978
+ info("cleaning up...");
1979
1975
- if (statsd.collection_threads) {
1976
- int i;
1977
- for (i = 0; i < statsd.threads; i++) {
1980
+ if (statsd.collection_threads && statsd.collection_threads_status) {
1981
+ int i;
1982
+ for (i = 0; i < statsd.threads; i++) {
1983
+ if(statsd.collection_threads_status[i]) {
1984
info("STATSD: stopping data collection thread %d...", i + 1);
1985
netdata_thread_cancel(statsd.collection_threads[i]);
1986
}
1987
+ else {
1988
+ info("STATSD: data collection thread %d found stopped.", i + 1);
1989
+ }
1990
}
1991
+ }
1992
1983
- info("STATSD: closing sockets...");
1984
- listen_sockets_close(&statsd.sockets);
1993
+ info("STATSD: closing sockets...");
1994
+ listen_sockets_close(&statsd.sockets);
1995
1986
- info("STATSD: cleanup completed.");
1987
- static_thread->enabled = 0;
1988
- }
1996
+ info("STATSD: cleanup completed.");
1997
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
1998
}
1999
2000
void *statsd_main(void *ptr) {
@@ -2081,12 +2090,14 @@ void *statsd_main(void *ptr) {
2090
goto cleanup;
2091
}
2092
2093
+ statsd.collection_threads_status = callocz((size_t)statsd.threads, sizeof(int));
2094
statsd.collection_threads = callocz((size_t)statsd.threads, sizeof(netdata_thread_t));
2095
+
2096
int i;
2097
for(i = 0; i < statsd.threads ;i++) {
2098
char tag[NETDATA_THREAD_TAG_MAX + 1];
2099
snprintfz(tag, NETDATA_THREAD_TAG_MAX, "STATSD_COLLECTOR[%d]", i + 1);
2089
- netdata_thread_create(&statsd.collection_threads[i], tag, NETDATA_THREAD_OPTION_DEFAULT, statsd_collector_thread, &i);
2100
+ netdata_thread_create(&statsd.collection_threads[i], tag, NETDATA_THREAD_OPTION_DEFAULT, statsd_collector_thread, &statsd.collection_threads_status[i]);
2101
}
2102
2103
// ----------------------------------------------------------------------------------------------------------------
src/sys_fs_cgroup.c
+4
-4
@@ -2676,11 +2676,11 @@ void update_cgroup_charts(int update_every) {
2676
2677
static void cgroup_main_cleanup(void *ptr) {
2678
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
2679
- if(static_thread->enabled) {
2680
- info("cleaning up...");
2679
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
2680
2682
- static_thread->enabled = 0;
2683
- }
2681
+ info("cleaning up...");
2682
+
2683
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
2684
}
2685
2686
void *cgroups_main(void *ptr) {
src/web_server.c
+62
-51
@@ -211,6 +211,8 @@ static void web_client_cache_destroy(void) {
211
web_client_cache_verify(1);
212
#endif
213
214
+ netdata_thread_disable_cancelability();
215
+
216
struct web_client *w, *t;
217
218
w = web_clients_cache.used;
@@ -230,6 +232,8 @@ static void web_client_cache_destroy(void) {
232
}
233
web_clients_cache.avail = NULL;
234
web_clients_cache.avail_count = 0;
235
+
236
+ netdata_thread_enable_cancelability();
237
}
238
239
static struct web_client *web_client_get_from_cache_or_allocate() {
@@ -242,6 +246,8 @@ static struct web_client *web_client_get_from_cache_or_allocate() {
246
error("Oops! wrong thread accessing the cache. Expected %d, found %d", (int)web_clients_cache.pid, (int)gettid());
247
#endif
248
249
+ netdata_thread_disable_cancelability();
250
+
251
struct web_client *w = web_clients_cache.avail;
252
253
if(w) {
@@ -269,6 +275,9 @@ static struct web_client *web_client_get_from_cache_or_allocate() {
275
// initialize it
276
w->id = web_client_connected();
277
w->mode = WEB_CLIENT_MODE_NORMAL;
278
+
279
+ netdata_thread_enable_cancelability();
280
+
281
return w;
282
}
283
@@ -283,12 +292,16 @@ static void web_client_release(struct web_client *w) {
292
293
debug(D_WEB_CLIENT_ACCESS, "%llu: Closing web client from %s port %s.", w->id, w->client_ip, w->client_port);
294
295
+ log_connection(w, "DISCONNECTED");
296
web_client_request_done(w);
297
web_client_disconnected();
298
299
+ netdata_thread_disable_cancelability();
300
+
301
if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
302
if (w->ifd != -1) close(w->ifd);
303
if (w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
304
+ w->ifd = w->ofd = -1;
305
}
306
307
// unlink it from the used
@@ -309,6 +322,8 @@ static void web_client_release(struct web_client *w) {
322
web_clients_cache.avail = w;
323
web_clients_cache.avail_count++;
324
}
325
+
326
+ netdata_thread_enable_cancelability();
327
}
328
329
@@ -580,25 +595,24 @@ static struct pollfd *socket_listen_main_multi_threaded_fds = NULL;
595
596
static void socket_listen_main_multi_threaded_cleanup(void *data) {
597
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)data;
583
- if(static_thread->enabled) {
584
- static_thread->enabled = 0;
598
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
599
586
- info("cleaning up...");
600
+ info("cleaning up...");
601
588
- info("releasing allocated memory...");
589
- freez(socket_listen_main_multi_threaded_fds);
602
+ info("releasing allocated memory...");
603
+ freez(socket_listen_main_multi_threaded_fds);
604
591
- info("closing all sockets...");
592
- listen_sockets_close(&api_sockets);
605
+ info("closing all sockets...");
606
+ listen_sockets_close(&api_sockets);
607
594
- info("stopping all running web server threads...");
595
- web_client_multi_threaded_web_server_stop_all_threads();
608
+ info("stopping all running web server threads...");
609
+ web_client_multi_threaded_web_server_stop_all_threads();
610
597
- info("freeing web clients cache...");
598
- web_client_cache_destroy();
611
+ info("freeing web clients cache...");
612
+ web_client_cache_destroy();
613
600
- info("cleanup completed.");
601
- }
614
+ info("cleanup completed.");
615
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
616
}
617
618
#define CLEANUP_EVERY_EVENTS 60
@@ -734,20 +748,16 @@ static inline int single_threaded_unlink_client(struct web_client *w, fd_set *if
748
749
static void socket_listen_main_single_threaded_cleanup(void *data) {
750
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)data;
737
- if(static_thread->enabled) {
738
- static_thread->enabled = 0;
739
-
740
- info("cleaning up...");
751
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
752
742
- info("closing all sockets...");
743
- listen_sockets_close(&api_sockets);
753
+ info("closing all sockets...");
754
+ listen_sockets_close(&api_sockets);
755
745
- info("freeing web clients cache...");
746
- web_client_cache_destroy();
756
+ info("freeing web clients cache...");
757
+ web_client_cache_destroy();
758
748
- info("cleanup completed.");
749
- debug(D_WEB_CLIENT, "LISTENER: exit!");
750
- }
759
+ info("cleanup completed.");
760
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
761
}
762
763
void *socket_listen_main_single_threaded(void *ptr) {
@@ -1141,41 +1151,42 @@ void *socket_listen_main_static_threaded_worker(void *ptr) {
1151
1152
static void socket_listen_main_static_threaded_cleanup(void *ptr) {
1153
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
1144
- if(static_thread->enabled) {
1145
- int i, found = 0, max = 2 * USEC_PER_SEC, step = 50000;
1154
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
1155
1147
- // we start from 1, - 0 is self
1148
- for(i = 1; i < static_threaded_workers_count; i++) {
1149
- if(static_workers_private_data[i].running) {
1150
- found++;
1151
- info("stopping worker %d", i + 1);
1152
- netdata_thread_cancel(static_workers_private_data[i].thread);
1153
- }
1154
- else
1155
- info("found stopped worker %d", i + 1);
1156
+ int i, found = 0, max = 2 * USEC_PER_SEC, step = 50000;
1157
+
1158
+ // we start from 1, - 0 is self
1159
+ for(i = 1; i < static_threaded_workers_count; i++) {
1160
+ if(static_workers_private_data[i].running) {
1161
+ found++;
1162
+ info("stopping worker %d", i + 1);
1163
+ netdata_thread_cancel(static_workers_private_data[i].thread);
1164
}
1165
+ else
1166
+ info("found stopped worker %d", i + 1);
1167
+ }
1168
1158
- while(found && max > 0) {
1159
- max -= step;
1160
- info("Waiting %d static web threads to finish...", found);
1161
- sleep_usec(step);
1162
- found = 0;
1169
+ while(found && max > 0) {
1170
+ max -= step;
1171
+ info("Waiting %d static web threads to finish...", found);
1172
+ sleep_usec(step);
1173
+ found = 0;
1174
1164
- // we start from 1, - 0 is self
1165
- for(i = 1; i < static_threaded_workers_count; i++) {
1166
- if (static_workers_private_data[i].running)
1167
- found++;
1168
- }
1175
+ // we start from 1, - 0 is self
1176
+ for(i = 1; i < static_threaded_workers_count; i++) {
1177
+ if (static_workers_private_data[i].running)
1178
+ found++;
1179
}
1180
+ }
1181
1171
- if(found)
1172
- error("%d static web threads are taking too long to finish. Giving up.", found);
1182
+ if(found)
1183
+ error("%d static web threads are taking too long to finish. Giving up.", found);
1184
1174
- info("closing all web server sockets...");
1175
- listen_sockets_close(&api_sockets);
1185
+ info("closing all web server sockets...");
1186
+ listen_sockets_close(&api_sockets);
1187
1177
- static_thread->enabled = 0;
1178
- }
1188
+ info("all static web threads stopped.");
1189
+ static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
1190
}
1191
1192
void *socket_listen_main_static_threaded(void *ptr) {