added caching of allocations among web server connections
Costa Tsaousis (ktsaou) committed
Jan 9, 2018 at 23:40 UTC
350e3b4ebd0475d4749fb77c804f446d14655525
17 files changed
+262
-183
src/backends.c
+2
-2
@@ -501,9 +501,9 @@ inline uint32_t backend_parse_data_source(const char *source, uint32_t mode) {
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
- static_thread->enabled = 0;
505
-
504
info("cleaning up...");
505
+
506
+ static_thread->enabled = 0;
507
}
508
}
509
src/health.c
+2
-2
@@ -341,9 +341,9 @@ static inline int check_if_resumed_from_suspention(void) {
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
- static_thread->enabled = 0;
345
-
344
info("cleaning up...");
345
+
346
+ static_thread->enabled = 0;
347
}
348
}
349
src/main.c
+1
-1
@@ -196,7 +196,7 @@ int killpid(pid_t pid, int sig)
196
void cancel_main_threads() {
197
error_log_limit_unlimited();
198
199
- int i, found = 0, max = 1 * USEC_PER_SEC, step = 100000;
199
+ int i, found = 0, max = 5 * USEC_PER_SEC, step = 100000;
200
for (i = 0; static_threads[i].name != NULL ; i++) {
201
if(static_threads[i].enabled) {
202
info("EXIT: Stopping master thread: %s", static_threads[i].name);
src/plugin_checks.c
+2
-2
@@ -5,9 +5,9 @@
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
- static_thread->enabled = 0;
9
-
8
info("cleaning up...");
9
+
10
+ static_thread->enabled = 0;
11
}
12
}
13
src/plugin_freebsd.c
+2
-2
@@ -69,9 +69,9 @@ static struct freebsd_module {
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
- static_thread->enabled = 0;
73
-
72
info("cleaning up...");
73
+
74
+ static_thread->enabled = 0;
75
}
76
}
77
src/plugin_idlejitter.c
+2
-2
@@ -5,9 +5,9 @@
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
- static_thread->enabled = 0;
9
-
8
info("cleaning up...");
9
+
10
+ static_thread->enabled = 0;
11
}
12
}
13
src/plugin_macos.c
+2
-2
@@ -3,9 +3,9 @@
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
- static_thread->enabled = 0;
7
-
6
info("cleaning up...");
7
+
8
+ static_thread->enabled = 0;
9
}
10
}
11
src/plugin_nfacct.c
+2
-2
@@ -754,8 +754,6 @@ static void nfacct_send_metrics() {
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("cleaning up...");
758
759
#ifdef DO_NFACCT
@@ -765,6 +763,8 @@ static void nfacct_main_cleanup(void *ptr) {
763
#ifdef DO_NFSTAT
764
nfstat_cleanup();
765
#endif
766
+
767
+ static_thread->enabled = 0;
768
}
769
}
770
src/plugin_proc.c
+2
-2
@@ -67,9 +67,9 @@ static struct proc_module {
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
- static_thread->enabled = 0;
71
-
70
info("cleaning up...");
71
+
72
+ static_thread->enabled = 0;
73
}
74
}
75
src/plugin_proc_diskspace.c
+2
-2
@@ -331,9 +331,9 @@ static inline void do_disk_space_stats(struct mountinfo *mi, int update_every) {
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
- static_thread->enabled = 0;
335
-
334
info("cleaning up...");
335
+
336
+ static_thread->enabled = 0;
337
}
338
}
339
src/plugin_tc.c
+2
-2
@@ -833,8 +833,6 @@ static pid_t tc_child_pid = 0;
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
- static_thread->enabled = 0;
837
-
836
info("cleaning up...");
837
838
if(tc_child_pid) {
@@ -849,6 +847,8 @@ static void tc_main_cleanup(void *ptr) {
847
848
tc_child_pid = 0;
849
}
850
+
851
+ static_thread->enabled = 0;
852
}
853
}
854
src/plugins_d.c
+1
-2
@@ -569,8 +569,6 @@ void *pluginsd_worker_thread(void *arg) {
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
- static_thread->enabled = 0;
573
-
572
info("cleaning up...");
573
574
struct plugind *cd;
@@ -582,6 +580,7 @@ static void pluginsd_main_cleanup(void *data) {
580
}
581
582
info("cleanup completed.");
583
+ static_thread->enabled = 0;
584
}
585
}
586
src/statsd.c
+1
-2
@@ -1970,8 +1970,6 @@ static int statsd_listen_sockets_setup(void) {
1970
static void statsd_main_cleanup(void *data) {
1971
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)data;
1972
if(static_thread->enabled) {
1973
- static_thread->enabled = 0;
1974
-
1973
info("cleaning up...");
1974
1975
if (statsd.collection_threads) {
@@ -1986,6 +1984,7 @@ static void statsd_main_cleanup(void *data) {
1984
listen_sockets_close(&statsd.sockets);
1985
1986
info("STATSD: cleanup completed.");
1987
+ static_thread->enabled = 0;
1988
}
1989
}
1990
src/sys_fs_cgroup.c
+2
-2
@@ -2677,9 +2677,9 @@ void update_cgroup_charts(int update_every) {
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
- static_thread->enabled = 0;
2681
-
2680
info("cleaning up...");
2681
+
2682
+ static_thread->enabled = 0;
2683
}
2684
}
2685
src/web_client.c
+165
-41
@@ -92,21 +92,167 @@ static void web_client_initialize_connection(struct web_client *w) {
92
93
web_client_update_acl_matches(w);
94
95
- w->response.data = buffer_create(INITIAL_WEB_DATA_LENGTH);
96
- w->response.header = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
97
- w->response.header_output = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
98
- w->origin[0] = '*';
95
+ w->origin[0] = '*'; w->origin[1] = '\0';
96
web_client_enable_wait_receive(w);
97
98
log_connection(w, "CONNECTED");
99
}
100
104
-struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port) {
101
+struct clients_cache {
102
+ struct web_client *used;
103
+ size_t used_count;
104
+
105
+ struct web_client *avail;
106
+ size_t avail_count;
107
+};
108
+
109
+__thread struct clients_cache web_clients_cache = {
110
+ .used = 0,
111
+ .used_count = 0,
112
+
113
+ .avail = 0,
114
+ .avail_count = 0
115
+};
116
+
117
+static void web_client_free(struct web_client *w) {
118
+ buffer_free(w->response.header_output);
119
+ buffer_free(w->response.header);
120
+ buffer_free(w->response.data);
121
+ freez(w);
122
+}
123
+
124
+void web_client_cache_destroy(void) {
125
+ struct web_client *w, *t;
126
+
127
+ w = web_clients_cache.used;
128
+ while(w) {
129
+ t = w;
130
+ w = w->next;
131
+ web_client_free(t);
132
+ }
133
+ web_clients_cache.used = NULL;
134
+ web_clients_cache.used_count = 0;
135
+
136
+ w = web_clients_cache.avail;
137
+ while(w) {
138
+ t = w;
139
+ w = w->next;
140
+ web_client_free(t);
141
+ }
142
+ web_clients_cache.avail = NULL;
143
+ web_clients_cache.avail_count = 0;
144
+}
145
+
146
+void web_client_multi_threaded_web_server_stop_all_threads(void) {
147
struct web_client *w;
148
107
- w = callocz(1, sizeof(struct web_client));
149
+ int found = 1, max = 2 * USEC_PER_SEC, step = 50000;
150
+ for(w = web_clients_cache.used; w ; w = w->next) {
151
+ if(w->running) {
152
+ found++;
153
+ info("stopping web client %s, id %llu", w->client_ip, w->id);
154
+ netdata_thread_cancel(w->thread);
155
+ }
156
+ }
157
+
158
+ while(found && max > 0) {
159
+ max -= step;
160
+ info("Waiting %d web threads to finish...", found);
161
+ sleep_usec(step);
162
+ found = 0;
163
+ for(w = web_clients_cache.used; w ; w = w->next)
164
+ if(w->running) found++;
165
+ }
166
+
167
+ if(found)
168
+ error("%d web threads are taking too long to finish. Giving up.", found);
169
+}
170
+
171
+static void web_client_return_to_cache_or_free(struct web_client *w) {
172
+ // unlink it from the used;
173
+ if (w == web_clients_cache.used) web_clients_cache.used = w->next;
174
+ if(w->prev) w->prev->next = w->next;
175
+ if(w->next) w->next->prev = w->prev;
176
+ web_clients_cache.used_count--;
177
+
178
+ if(web_clients_cache.avail_count > 100) {
179
+ // we have too many of them - free it
180
+ web_client_free(w);
181
+ }
182
+ else {
183
+ // link it to the avail
184
+ if (web_clients_cache.avail) web_clients_cache.avail->prev = w;
185
+ w->next = web_clients_cache.avail;
186
+ w->prev = NULL;
187
+ web_clients_cache.avail = w;
188
+ web_clients_cache.avail_count++;
189
+ }
190
+}
191
+
192
+static struct web_client *web_client_get_from_cache_or_allocate() {
193
+ struct web_client *w = web_clients_cache.avail;
194
+
195
+ if(w) {
196
+ // unlink it from avail
197
+ if (w == web_clients_cache.avail) web_clients_cache.avail = w->next;
198
+ if(w->prev) w->prev->next = w->next;
199
+ if(w->next) w->next->prev = w->prev;
200
+ web_clients_cache.avail_count--;
201
+
202
+ // zero everything about it - but keep the buffers
203
+
204
+ BUFFER *b1 = w->response.data;
205
+ BUFFER *b2 = w->response.header;
206
+ BUFFER *b3 = w->response.header_output;
207
+
208
+ buffer_flush(b1);
209
+ buffer_flush(b2);
210
+ buffer_flush(b3);
211
+
212
+ memset(w, 0, sizeof(struct web_client));
213
+
214
+ w->response.data = b1;
215
+ w->response.header = b2;
216
+ w->response.header_output = b3;
217
+ }
218
+ else {
219
+ w = callocz(1, sizeof(struct web_client));
220
+ w->response.data = buffer_create(INITIAL_WEB_DATA_LENGTH);
221
+ w->response.header = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
222
+ w->response.header_output = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
223
+ }
224
+
225
+ // link it to used web clients
226
+ if (web_clients_cache.used) web_clients_cache.used->prev = w;
227
+ w->next = web_clients_cache.used;
228
+ w->prev = NULL;
229
+ web_clients_cache.used = w;
230
+ web_clients_cache.used_count++;
231
+
232
+ // initialize it
233
w->id = web_client_connected();
234
w->mode = WEB_CLIENT_MODE_NORMAL;
235
+ return w;
236
+}
237
+
238
+void web_client_release(struct web_client *w) {
239
+ debug(D_WEB_CLIENT_ACCESS, "%llu: Closing web client from %s port %s.", w->id, w->client_ip, w->client_port);
240
+
241
+ web_client_request_done(w);
242
+ web_client_disconnected();
243
+
244
+ if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
245
+ if (w->ifd != -1) close(w->ifd);
246
+ if (w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
247
+ }
248
+
249
+ web_client_return_to_cache_or_free(w);
250
+}
251
+
252
+struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port) {
253
+ struct web_client *w;
254
+
255
+ w = web_client_get_from_cache_or_allocate();
256
w->ifd = w->ofd = fd;
257
258
strncpyz(w->client_ip, client_ip, sizeof(w->client_ip) - 1);
@@ -122,10 +268,7 @@ struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const
268
struct web_client *web_client_create_on_listenfd(int listener) {
269
struct web_client *w;
270
125
- w = callocz(1, sizeof(struct web_client));
126
- w->id = web_client_connected();
127
- w->mode = WEB_CLIENT_MODE_NORMAL;
128
-
271
+ w = web_client_get_from_cache_or_allocate();
272
w->ifd = w->ofd = accept_socket(listener, SOCK_NONBLOCK, w->client_ip, sizeof(w->client_ip), w->client_port, sizeof(w->client_port), web_allow_connections_from);
273
274
if(unlikely(!*w->client_ip)) strcpy(w->client_ip, "-");
@@ -139,8 +282,7 @@ struct web_client *web_client_create_on_listenfd(int listener) {
282
error("%llu: Failed to accept new incoming connection.", w->id);
283
}
284
142
- freez(w);
143
- web_client_disconnected();
285
+ web_client_release(w);
286
return NULL;
287
}
288
@@ -148,7 +290,7 @@ struct web_client *web_client_create_on_listenfd(int listener) {
290
return(w);
291
}
292
151
-void web_client_reset(struct web_client *w) {
293
+void web_client_request_done(struct web_client *w) {
294
web_client_uncrock_socket(w);
295
296
debug(D_WEB_CLIENT, "%llu: Resetting client.", w->id);
@@ -273,29 +415,6 @@ void web_client_reset(struct web_client *w) {
415
#endif // NETDATA_WITH_ZLIB
416
}
417
276
-void web_client_free(struct web_client *w) {
277
- debug(D_WEB_CLIENT_ACCESS, "%llu: Closing web client from %s port %s.", w->id, w->client_ip, w->client_port);
278
-
279
- web_client_reset(w);
280
-
281
- buffer_free(w->response.header_output);
282
- w->response.header_output = NULL;
283
-
284
- buffer_free(w->response.header);
285
- w->response.header = NULL;
286
-
287
- buffer_free(w->response.data);
288
- w->response.data = NULL;
289
-
290
- if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
291
- if (w->ifd != -1) close(w->ifd);
292
- if (w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
293
- }
294
-
295
- freez(w);
296
- web_client_disconnected();
297
-}
298
-
418
uid_t web_files_uid(void) {
419
static char *web_owner = NULL;
420
static uid_t owner_uid = 0;
@@ -1453,7 +1572,7 @@ void web_client_process_request(struct web_client *w) {
1572
if(len != w->response.data->rbytes)
1573
error("%llu: sendfile() should copy %ld bytes, but copied %ld. Falling back to manual copy.", w->id, w->response.data->rbytes, len);
1574
else
1456
- web_client_reset(w);
1575
+ web_client_request_done(w);
1576
}
1577
*/
1578
}
@@ -1571,7 +1690,7 @@ ssize_t web_client_send_deflate(struct web_client *w)
1690
}
1691
1692
// reset the client
1574
- web_client_reset(w);
1693
+ web_client_request_done(w);
1694
debug(D_WEB_CLIENT, "%llu: Done sending all data on socket.", w->id);
1695
return t;
1696
}
@@ -1611,7 +1730,7 @@ ssize_t web_client_send_deflate(struct web_client *w)
1730
// compress
1731
if(deflate(&w->response.zstream, flush) == Z_STREAM_ERROR) {
1732
error("%llu: Compression failed. Closing down client.", w->id);
1614
- web_client_reset(w);
1733
+ web_client_request_done(w);
1734
return(-1);
1735
}
1736
@@ -1682,7 +1801,7 @@ ssize_t web_client_send(struct web_client *w) {
1801
return 0;
1802
}
1803
1685
- web_client_reset(w);
1804
+ web_client_request_done(w);
1805
debug(D_WEB_CLIENT, "%llu: Done sending all data on socket. Waiting for next request on the same socket.", w->id);
1806
return 0;
1807
}
@@ -1794,15 +1913,20 @@ ssize_t web_client_receive(struct web_client *w)
1913
1914
static void web_client_main_cleanup(void *ptr) {
1915
struct web_client *w = ptr;
1916
+
1917
if(!web_client_check_obsolete(w)) {
1918
WEB_CLIENT_IS_OBSOLETE(w);
1919
}
1920
+
1921
+ w->running = 0;
1922
}
1923
1924
void *web_client_main(void *ptr) {
1925
netdata_thread_cleanup_push(web_client_main_cleanup, ptr);
1926
1927
struct web_client *w = ptr;
1928
+ w->running = 1;
1929
+
1930
struct pollfd fds[2], *ifd, *ofd;
1931
int retval, timeout;
1932
nfds_t fdmax = 0;
@@ -1915,7 +2039,7 @@ void *web_client_main(void *ptr) {
2039
if(w->mode != WEB_CLIENT_MODE_STREAM)
2040
log_connection(w, "DISCONNECTED");
2041
1918
- web_client_reset(w);
2042
+ web_client_request_done(w);
2043
2044
debug(D_WEB_CLIENT, "%llu: done...", w->id);
2045
src/web_client.h
+13
-7
@@ -155,14 +155,17 @@ struct web_client {
155
size_t stats_received_bytes;
156
size_t stats_sent_bytes;
157
158
- // STATIC-THREADED WEB SERVER MEMBERS
159
- size_t pollinfo_slot; // POLLINFO slot of the web client
160
- size_t pollinfo_filecopy_slot; // POLLINFO slot of the file read
158
+ // cache of web_client allocations
159
+ struct web_client *prev; // maintain a linked list of web clients
160
+ struct web_client *next; // for the web servers that need it
161
162
// MULTI-THREADED WEB SERVER MEMBERS
163
netdata_thread_t thread; // the thread servicing this client
164
- struct web_client *prev; // maintain a linked list of web clients
165
- struct web_client *next; // for the web servers that need it
164
+ volatile int running; // 1 when the thread runs, 0 otherwise
165
+
166
+ // STATIC-THREADED WEB SERVER MEMBERS
167
+ size_t pollinfo_slot; // POLLINFO slot of the web client
168
+ size_t pollinfo_filecopy_slot; // POLLINFO slot of the file read
169
};
170
171
extern SIMPLE_PATTERN *web_allow_connections_from;
@@ -179,14 +182,14 @@ extern int web_client_permission_denied(struct web_client *w);
182
183
extern struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port);
184
extern struct web_client *web_client_create_on_listenfd(int listener);
182
-extern void web_client_free(struct web_client *w);
185
+extern void web_client_release(struct web_client *w);
186
187
extern ssize_t web_client_send(struct web_client *w);
188
extern ssize_t web_client_receive(struct web_client *w);
189
extern ssize_t web_client_read_file(struct web_client *w);
190
191
extern void web_client_process_request(struct web_client *w);
189
-extern void web_client_reset(struct web_client *w);
192
+extern void web_client_request_done(struct web_client *w);
193
194
extern void *web_client_main(void *ptr);
195
@@ -197,4 +200,7 @@ extern void buffer_data_options2string(BUFFER *wb, uint32_t options);
200
201
extern int mysendfile(struct web_client *w, char *filename);
202
203
+extern void web_client_multi_threaded_web_server_stop_all_threads(void);
204
+extern void web_client_cache_destroy(void);
205
+
206
#endif
src/web_server.c
+59
-108
@@ -11,40 +11,6 @@ WEB_SERVER_MODE web_server_mode = WEB_SERVER_MODE_STATIC_THREADED;
11
12
// --------------------------------------------------------------------------------------
13
14
-#ifdef NETDATA_INTERNAL_CHECKS
15
-static void log_allocations(void)
16
-{
17
-#ifdef HAVE_C_MALLINFO
18
- static int heap = 0, used = 0, mmap = 0;
19
-
20
- struct mallinfo mi;
21
-
22
- mi = mallinfo();
23
- if(mi.uordblks > used) {
24
- info("Allocated memory: used %d KB (+%d B), mmap %d KB (+%d B), heap %d KB (+%d B).",
25
- mi.uordblks / 1024,
26
- mi.uordblks - used,
27
- mi.hblkhd / 1024,
28
- mi.hblkhd - mmap,
29
- mi.arena / 1024,
30
- mi.arena - heap);
31
-
32
- used = mi.uordblks;
33
- heap = mi.arena;
34
- mmap = mi.hblkhd;
35
- }
36
-#else /* ! HAVE_C_MALLINFO */
37
- ;
38
-#endif /* ! HAVE_C_MALLINFO */
39
-
40
-#ifdef has_jemalloc
41
- malloc_stats_print(NULL, NULL, NULL);
42
-#endif
43
-}
44
-#endif /* NETDATA_INTERNAL_CHECKS */
45
-
46
-// --------------------------------------------------------------------------------------
47
-
14
WEB_SERVER_MODE web_server_mode_id(const char *mode) {
15
if(!strcmp(mode, "none"))
16
return WEB_SERVER_MODE_NONE;
@@ -87,50 +53,11 @@ int api_listen_sockets_setup(void) {
53
// --------------------------------------------------------------------------------------
54
// the main socket listener - MULTI-THREADED
55
90
-// maintain a linked list of web clients - for the web servers that need it
91
-static struct web_client *web_clients = NULL;
92
-
93
-static void web_client_link(struct web_client *w) {
94
- if (web_clients) web_clients->prev = w;
95
- w->next = web_clients;
96
- web_clients = w;
97
-}
98
-
99
-static struct web_client *web_client_unlink(struct web_client *w) {
100
- struct web_client *n = w->next;
101
- if (w == web_clients) web_clients = n;
102
-
103
- if(w->prev) w->prev->next = w->next;
104
- if(w->next) w->next->prev = w->prev;
105
-
106
- return n;
107
-}
108
-
109
-static inline void multi_threaded_cleanup_web_clients(void) {
110
- struct web_client *w;
111
-
112
- for (w = web_clients; w;) {
113
- if (web_client_check_obsolete(w)) {
114
- debug(D_WEB_CLIENT, "%llu: Removing client.", w->id);
115
- struct web_client *t = web_client_unlink(w);
116
- web_client_free(w);
117
- w = t;
118
-
119
-#ifdef NETDATA_INTERNAL_CHECKS
120
- log_allocations();
121
-#endif
122
- }
123
- else w = w->next;
124
- }
125
-}
126
-
56
// 1. it accepts new incoming requests on our port
57
// 2. creates a new web_client for each connection received
58
// 3. spawns a new netdata_thread to serve the client (this is optimal for keep-alive clients)
59
// 4. cleans up old web_clients that their netdata_threads have been exited
60
132
-#define CLEANUP_EVERY_EVENTS 100
133
-
61
static struct pollfd *socket_listen_main_multi_threaded_fds = NULL;
62
63
static void socket_listen_main_multi_threaded_cleanup(void *data) {
@@ -146,16 +73,13 @@ static void socket_listen_main_multi_threaded_cleanup(void *data) {
73
info("closing all sockets...");
74
listen_sockets_close(&api_sockets);
75
149
- info("cleanup completed.");
150
- }
76
+ info("stopping all running web server threads...");
77
+ web_client_multi_threaded_web_server_stop_all_threads();
78
152
- struct web_client *w;
153
- for(w = web_clients; w ; w = w->next) {
154
- if(!web_client_check_obsolete(w)) {
155
- WEB_CLIENT_IS_OBSOLETE(w);
156
- info("Stopping web client %s, id %llu", w->client_ip, w->id);
157
- netdata_thread_cancel(w->thread);
158
- }
79
+ info("freeing web clients cache...");
80
+ web_client_cache_destroy();
81
+
82
+ info("cleanup completed.");
83
}
84
}
85
@@ -166,7 +90,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
90
web_server_is_multithreaded = 1;
91
92
struct web_client *w;
169
- int retval, counter = 0;
93
+ int retval;
94
95
if(!api_sockets.opened)
96
fatal("LISTENER: No sockets to listen to.");
@@ -194,9 +118,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
118
continue;
119
}
120
else if(unlikely(!retval)) {
197
- debug(D_WEB_CLIENT, "LISTENER: select() timeout.");
198
- counter = 0;
199
- multi_threaded_cleanup_web_clients();
121
+ debug(D_WEB_CLIENT, "LISTENER: poll() timeout.");
122
continue;
123
}
124
@@ -212,7 +134,6 @@ void *socket_listen_main_multi_threaded(void *ptr) {
134
// no need for error log - web_client_create_on_listenfd already logged the error
135
continue;
136
}
215
- web_client_link(w);
137
138
if(api_sockets.fds_families[i] == AF_UNIX)
139
web_client_set_unix(w);
@@ -226,13 +147,6 @@ void *socket_listen_main_multi_threaded(void *ptr) {
147
WEB_CLIENT_IS_OBSOLETE(w);
148
}
149
}
229
-
230
- // cleanup unused clients
231
- counter++;
232
- if(counter >= CLEANUP_EVERY_EVENTS) {
233
- counter = 0;
234
- multi_threaded_cleanup_web_clients();
235
- }
150
}
151
152
netdata_thread_cleanup_pop(1);
@@ -245,8 +159,10 @@ void *socket_listen_main_multi_threaded(void *ptr) {
159
struct web_client *single_threaded_clients[FD_SETSIZE];
160
161
static inline int single_threaded_link_client(struct web_client *w, fd_set *ifds, fd_set *ofds, fd_set *efds, int *max) {
248
- if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w))))
162
+ if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w)))) {
163
+ // error("refusing to link obsolete/dead client");
164
return 1;
165
+ }
166
167
if(unlikely(w->ifd < 0 || w->ifd >= (int)FD_SETSIZE || w->ofd < 0 || w->ofd >= (int)FD_SETSIZE)) {
168
error("%llu: invalid file descriptor, ifd = %d, ofd = %d (required 0 <= fd < FD_SETSIZE (%d)", w->id, w->ifd, w->ofd, (int)FD_SETSIZE);
@@ -280,8 +196,10 @@ static inline int single_threaded_unlink_client(struct web_client *w, fd_set *if
196
single_threaded_clients[w->ifd] = NULL;
197
single_threaded_clients[w->ofd] = NULL;
198
283
- if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w))))
199
+ if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w)))) {
200
+ // error("unlinked client is obsolete/dead");
201
return 1;
202
+ }
203
204
return 0;
205
}
@@ -296,6 +214,9 @@ static void socket_listen_main_single_threaded_cleanup(void *data) {
214
info("closing all sockets...");
215
listen_sockets_close(&api_sockets);
216
217
+ info("freeing web clients cache...");
218
+ web_client_cache_destroy();
219
+
220
info("cleanup completed.");
221
debug(D_WEB_CLIENT, "LISTENER: exit!");
222
}
@@ -362,7 +283,7 @@ void *socket_listen_main_single_threaded(void *ptr) {
283
web_client_set_tcp(w);
284
285
if (single_threaded_link_client(w, &ifds, &ofds, &ifds, &fdmax) != 0) {
365
- web_client_free(w);
286
+ web_client_release(w);
287
}
288
}
289
}
@@ -372,22 +293,27 @@ void *socket_listen_main_single_threaded(void *ptr) {
293
continue;
294
295
w = single_threaded_clients[i];
375
- if(unlikely(!w))
296
+ if(unlikely(!w)) {
297
+ // error("no client on slot %zu", i);
298
continue;
299
+ }
300
301
if(unlikely(single_threaded_unlink_client(w, &ifds, &ofds, &efds) != 0)) {
379
- web_client_free(w);
302
+ // error("failed to unlink client %zu", i);
303
+ web_client_release(w);
304
continue;
305
}
306
307
if (unlikely(FD_ISSET(w->ifd, &refds) || FD_ISSET(w->ofd, &refds))) {
384
- web_client_free(w);
308
+ // error("no input on client %zu", i);
309
+ web_client_release(w);
310
continue;
311
}
312
313
if (unlikely(web_client_has_wait_receive(w) && FD_ISSET(w->ifd, &rifds))) {
314
if (unlikely(web_client_receive(w) < 0)) {
390
- web_client_free(w);
315
+ // error("cannot read from client %zu", i);
316
+ web_client_release(w);
317
continue;
318
}
319
@@ -399,14 +325,16 @@ void *socket_listen_main_single_threaded(void *ptr) {
325
326
if (unlikely(web_client_has_wait_send(w) && FD_ISSET(w->ofd, &rofds))) {
327
if (unlikely(web_client_send(w) < 0)) {
328
+ // error("cannot send data to client %zu", i);
329
debug(D_WEB_CLIENT, "%llu: Cannot send data to client. Closing client.", w->id);
403
- web_client_free(w);
330
+ web_client_release(w);
331
continue;
332
}
333
}
334
335
if(unlikely(single_threaded_link_client(w, &ifds, &ofds, &efds, &fdmax) != 0)) {
409
- web_client_free(w);
336
+ // error("failed to link client %zu", i);
337
+ web_client_release(w);
338
}
339
}
340
}
@@ -472,7 +400,7 @@ static void web_werver_file_del_callback(POLLINFO *pi) {
400
401
if(unlikely(!w->pollinfo_slot)) {
402
debug(D_WEB_CLIENT, "%llu: CROSS WEB CLIENT CLEANUP (iFD %d, oFD %d)", w->id, pi->fd, w->ofd);
475
- web_client_free(w);
403
+ web_client_release(w);
404
}
405
}
406
@@ -563,7 +491,7 @@ static void web_server_del_callback(POLLINFO *pi) {
491
}
492
else {
493
debug(D_WEB_CLIENT, "%llu: CLOSING CLIENT FD %d", w->id, pi->fd);
566
- web_client_free(w);
494
+ web_client_release(w);
495
}
496
}
497
@@ -642,7 +570,9 @@ static int web_server_snd_callback(POLLINFO *pi, short int *events) {
570
571
static void socket_listen_main_static_threaded_worker_cleanup(void *ptr) {
572
worker_private = (struct web_server_static_threaded_worker *)ptr;
645
- worker_private->running = 0;
573
+
574
+ info("freeing local web clients cache...");
575
+ web_client_cache_destroy();
576
577
info("stopped after %zu connects, %zu disconnects (max concurrent %zu), %zu receptions and %zu sends",
578
worker_private->connected,
@@ -651,6 +581,8 @@ static void socket_listen_main_static_threaded_worker_cleanup(void *ptr) {
581
worker_private->receptions,
582
worker_private->sends
583
);
584
+
585
+ worker_private->running = 0;
586
}
587
588
void *socket_listen_main_static_threaded_worker(void *ptr) {
@@ -679,11 +611,12 @@ void *socket_listen_main_static_threaded_worker(void *ptr) {
611
static void socket_listen_main_static_threaded_cleanup(void *ptr) {
612
struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
613
if(static_thread->enabled) {
682
- static_thread->enabled = 0;
614
+ int i, found = 0, max = 2 * USEC_PER_SEC, step = 50000;
615
684
- int i;
616
+ // we start from 1, - 0 is self
617
for(i = 1; i < static_threaded_workers_count; i++) {
618
if(static_workers_private_data[i].running) {
619
+ found++;
620
info("stopping worker %d", i + 1);
621
netdata_thread_cancel(static_workers_private_data[i].thread);
622
}
@@ -691,8 +624,26 @@ static void socket_listen_main_static_threaded_cleanup(void *ptr) {
624
info("found stopped worker %d", i + 1);
625
}
626
627
+ while(found && max > 0) {
628
+ max -= step;
629
+ info("Waiting %d static web threads to finish...", found);
630
+ sleep_usec(step);
631
+ found = 0;
632
+
633
+ // we start from 1, - 0 is self
634
+ for(i = 1; i < static_threaded_workers_count; i++) {
635
+ if (static_workers_private_data[i].running)
636
+ found++;
637
+ }
638
+ }
639
+
640
+ if(found)
641
+ error("%d static web threads are taking too long to finish. Giving up.", found);
642
+
643
info("closing all web server sockets...");
644
listen_sockets_close(&api_sockets);
645
+
646
+ static_thread->enabled = 0;
647
}
648
}
649