/src/fluent-bit/lib/monkey/mk_server/mk_scheduler.c
Line | Count | Source |
1 | | /* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ |
2 | | |
3 | | /* Monkey HTTP Server |
4 | | * ================== |
5 | | * Copyright 2001-2017 Eduardo Silva <eduardo@monkey.io> |
6 | | * |
7 | | * Licensed under the Apache License, Version 2.0 (the "License"); |
8 | | * you may not use this file except in compliance with the License. |
9 | | * You may obtain a copy of the License at |
10 | | * |
11 | | * http://www.apache.org/licenses/LICENSE-2.0 |
12 | | * |
13 | | * Unless required by applicable law or agreed to in writing, software |
14 | | * distributed under the License is distributed on an "AS IS" BASIS, |
15 | | * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
16 | | * See the License for the specific language governing permissions and |
17 | | * limitations under the License. |
18 | | */ |
19 | | |
20 | | #include <monkey/monkey.h> |
21 | | #include <monkey/mk_info.h> |
22 | | #include <monkey/mk_core.h> |
23 | | #include <monkey/mk_vhost.h> |
24 | | #include <monkey/mk_scheduler.h> |
25 | | #include <monkey/mk_scheduler_tls.h> |
26 | | #include <monkey/mk_server.h> |
27 | | #include <monkey/mk_thread.h> |
28 | | #include <monkey/mk_cache.h> |
29 | | #include <monkey/mk_config.h> |
30 | | #include <monkey/mk_clock.h> |
31 | | #include <monkey/mk_plugin.h> |
32 | | #include <monkey/mk_utils.h> |
33 | | #include <monkey/mk_linuxtrace.h> |
34 | | #include <monkey/mk_server.h> |
35 | | #include <monkey/mk_plugin_stage.h> |
36 | | #include <monkey/mk_http_thread.h> |
37 | | #include <monkey/mk_tls_transport.h> |
38 | | |
39 | | #include <signal.h> |
40 | | |
41 | | #ifndef _WIN32 |
42 | | #include <sys/syscall.h> |
43 | | #endif |
44 | | |
45 | | extern struct mk_sched_handler mk_http_handler; |
46 | | extern struct mk_sched_handler mk_http2_handler; |
47 | | |
48 | | pthread_mutex_t mutex_worker_init = PTHREAD_MUTEX_INITIALIZER; |
49 | | pthread_mutex_t mutex_worker_exit = PTHREAD_MUTEX_INITIALIZER; |
50 | | |
51 | | /* |
52 | | * Returns the worker id which should take a new incomming connection, |
53 | | * it returns the worker id with less active connections. Just used |
54 | | * if config->scheduler_mode is MK_SCHEDULER_FAIR_BALANCING. |
55 | | */ |
56 | | static inline int _next_target(struct mk_server *server) |
57 | 0 | { |
58 | 0 | int i; |
59 | 0 | int target = 0; |
60 | 0 | unsigned long long tmp = 0, cur = 0; |
61 | 0 | struct mk_sched_ctx *ctx = server->sched_ctx; |
62 | 0 | struct mk_sched_worker *worker; |
63 | |
|
64 | 0 | cur = (ctx->workers[0].accepted_connections - |
65 | 0 | ctx->workers[0].closed_connections); |
66 | 0 | if (cur == 0) { |
67 | 0 | return 0; |
68 | 0 | } |
69 | | |
70 | | /* Finds the lowest load worker */ |
71 | 0 | for (i = 1; i < server->workers; i++) { |
72 | 0 | worker = &ctx->workers[i]; |
73 | 0 | tmp = worker->accepted_connections - worker->closed_connections; |
74 | 0 | if (tmp < cur) { |
75 | 0 | target = i; |
76 | 0 | cur = tmp; |
77 | |
|
78 | 0 | if (cur == 0) |
79 | 0 | break; |
80 | 0 | } |
81 | 0 | } |
82 | | |
83 | | /* |
84 | | * If sched_ctx->workers[target] worker is full then the whole server too, |
85 | | * because it has the lowest load. |
86 | | */ |
87 | 0 | if (mk_unlikely(server->server_capacity > 0 && |
88 | 0 | server->server_capacity <= cur)) { |
89 | 0 | MK_TRACE("Too many clients: %i", server->server_capacity); |
90 | | |
91 | | /* Instruct to close the connection anyways, we lie, it will die */ |
92 | 0 | return -1; |
93 | 0 | } |
94 | | |
95 | 0 | return target; |
96 | 0 | } |
97 | | |
98 | | struct mk_sched_worker *mk_sched_next_target(struct mk_server *server) |
99 | 0 | { |
100 | 0 | int t; |
101 | 0 | struct mk_sched_ctx *ctx = server->sched_ctx; |
102 | |
|
103 | 0 | t = _next_target(server); |
104 | 0 | if (mk_likely(t != -1)) { |
105 | 0 | return &ctx->workers[t]; |
106 | 0 | } |
107 | | |
108 | 0 | return NULL; |
109 | 0 | } |
110 | | |
111 | | /* |
112 | | * This function is invoked when the core triggers a MK_SCHED_SIGNAL_FREE_ALL |
113 | | * event through the signal channels, it means the server will stop working |
114 | | * so this is the last call to release all memory resources in use. Of course |
115 | | * this takes place in a thread context. |
116 | | */ |
117 | | void mk_sched_worker_free(struct mk_server *server) |
118 | 0 | { |
119 | 0 | int i; |
120 | 0 | pthread_t tid; |
121 | 0 | struct mk_sched_ctx *ctx = server->sched_ctx; |
122 | 0 | struct mk_sched_worker *worker = NULL; |
123 | |
|
124 | 0 | pthread_mutex_lock(&mutex_worker_exit); |
125 | | |
126 | | /* |
127 | | * Fix Me: needs to implement API to make plugins release |
128 | | * their resources first at WORKER LEVEL |
129 | | */ |
130 | | |
131 | | /* External */ |
132 | 0 | mk_plugin_exit_worker(); |
133 | 0 | mk_vhost_fdt_worker_exit(server); |
134 | 0 | mk_cache_worker_exit(); |
135 | | |
136 | | /* Scheduler stuff */ |
137 | 0 | tid = pthread_self(); |
138 | 0 | for (i = 0; i < server->workers; i++) { |
139 | 0 | worker = &ctx->workers[i]; |
140 | 0 | if (worker->tid == tid) { |
141 | 0 | break; |
142 | 0 | } |
143 | 0 | worker = NULL; |
144 | 0 | } |
145 | |
|
146 | 0 | mk_bug(!worker); |
147 | | |
148 | | /* FIXME!: there is nothing done here with the worker context */ |
149 | | |
150 | | /* Free master array (av queue & busy queue) */ |
151 | 0 | mk_mem_free(MK_TLS_GET(mk_tls_sched_cs)); |
152 | 0 | mk_mem_free(MK_TLS_GET(mk_tls_sched_cs_incomplete)); |
153 | 0 | mk_mem_free(MK_TLS_GET(mk_tls_sched_worker_notif)); |
154 | 0 | pthread_mutex_unlock(&mutex_worker_exit); |
155 | 0 | } |
156 | | |
157 | | struct mk_sched_handler *mk_sched_handler_cap(char cap) |
158 | 0 | { |
159 | 0 | if (cap == MK_CAP_HTTP) { |
160 | 0 | return &mk_http_handler; |
161 | 0 | } |
162 | | |
163 | | #ifdef MK_HAVE_HTTP2 |
164 | | else if (cap == MK_CAP_HTTP2) { |
165 | | return &mk_http2_handler; |
166 | | } |
167 | | #endif |
168 | | |
169 | 0 | return NULL; |
170 | 0 | } |
171 | | |
172 | | /* |
173 | | * Register a new client connection into the scheduler, this call takes place |
174 | | * inside the worker/thread context. |
175 | | */ |
176 | | struct mk_sched_conn *mk_sched_add_connection(int remote_fd, |
177 | | struct mk_server_listen *listener, |
178 | | struct mk_sched_worker *sched, |
179 | | struct mk_server *server) |
180 | 0 | { |
181 | 0 | int ret; |
182 | 0 | int size; |
183 | 0 | struct mk_sched_handler *handler; |
184 | 0 | struct mk_sched_conn *conn; |
185 | 0 | struct mk_event *event; |
186 | | |
187 | | /* Before to continue, we need to run plugin stage 10 */ |
188 | 0 | ret = mk_plugin_stage_run_10(remote_fd, server); |
189 | | |
190 | | /* Close connection, otherwise continue */ |
191 | 0 | if (ret == MK_PLUGIN_RET_CLOSE_CONX) { |
192 | 0 | listener->network->close(listener->network->plugin, remote_fd); |
193 | 0 | MK_LT_SCHED(remote_fd, "PLUGIN_CLOSE"); |
194 | 0 | return NULL; |
195 | 0 | } |
196 | | |
197 | 0 | handler = listener->protocol; |
198 | 0 | if (handler->sched_extra_size > 0) { |
199 | 0 | void *data; |
200 | 0 | size = (sizeof(struct mk_sched_conn) + handler->sched_extra_size); |
201 | 0 | data = mk_mem_alloc_z(size); |
202 | 0 | conn = (struct mk_sched_conn *) data; |
203 | 0 | } |
204 | 0 | else { |
205 | 0 | conn = mk_mem_alloc_z(sizeof(struct mk_sched_conn)); |
206 | 0 | } |
207 | |
|
208 | 0 | if (!conn) { |
209 | 0 | mk_err("[server] Could not register client"); |
210 | 0 | return NULL; |
211 | 0 | } |
212 | | |
213 | 0 | event = &conn->event; |
214 | 0 | event->fd = remote_fd; |
215 | 0 | event->type = MK_EVENT_CONNECTION; |
216 | 0 | event->mask = MK_EVENT_EMPTY; |
217 | 0 | event->status = MK_EVENT_NONE; |
218 | 0 | conn->arrive_time = server->clock_context->log_current_utime; |
219 | 0 | conn->protocol = handler; |
220 | 0 | conn->net = listener->network; |
221 | 0 | conn->is_timeout_on = MK_FALSE; |
222 | 0 | conn->server_listen = listener; |
223 | | |
224 | | /* Stream channel */ |
225 | 0 | conn->channel.type = MK_CHANNEL_SOCKET; /* channel type */ |
226 | 0 | conn->channel.fd = remote_fd; /* socket conn */ |
227 | 0 | conn->channel.io = conn->net; /* network layer */ |
228 | 0 | conn->channel.event = event; /* parent event ref */ |
229 | 0 | mk_list_init(&conn->channel.streams); |
230 | | |
231 | | /* |
232 | | * Register the connections into the timeout_queue: |
233 | | * |
234 | | * When a new connection arrives, we cannot assume it contains some data |
235 | | * to read, meaning the event loop may not get notifications and the protocol |
236 | | * handler will never be called. So in order to avoid DDoS we always register |
237 | | * this session in the timeout_queue for further lookup. |
238 | | * |
239 | | * The protocol handler is in charge to remove the session from the |
240 | | * timeout_queue. |
241 | | */ |
242 | 0 | mk_sched_conn_timeout_add(conn, sched); |
243 | | |
244 | | /* Linux trace message */ |
245 | 0 | MK_LT_SCHED(remote_fd, "REGISTERED"); |
246 | |
|
247 | 0 | return conn; |
248 | 0 | } |
249 | | |
250 | | static void mk_sched_thread_lists_init() |
251 | 0 | { |
252 | 0 | struct mk_list *sched_cs_incomplete; |
253 | | |
254 | | /* mk_tls_sched_cs_incomplete */ |
255 | 0 | sched_cs_incomplete = mk_mem_alloc(sizeof(struct mk_list)); |
256 | 0 | mk_list_init(sched_cs_incomplete); |
257 | 0 | MK_TLS_SET(mk_tls_sched_cs_incomplete, sched_cs_incomplete); |
258 | 0 | } |
259 | | |
260 | | /* Register thread information. The caller thread is the thread information's owner */ |
261 | | static int mk_sched_register_thread(struct mk_server *server) |
262 | 0 | { |
263 | 0 | struct mk_sched_ctx *ctx = server->sched_ctx; |
264 | 0 | struct mk_sched_worker *worker; |
265 | | |
266 | | /* |
267 | | * If this thread slept inside this section, some other thread may touch |
268 | | * server->worker_id. |
269 | | * So protect it with a mutex, only one thread may handle server->worker_id. |
270 | | * |
271 | | * Note : Let's use the platform agnostic atomics we implemented in cmetrics here |
272 | | * instead of a lock. |
273 | | */ |
274 | 0 | worker = &ctx->workers[server->worker_id]; |
275 | 0 | worker->idx = server->worker_id++; |
276 | 0 | worker->tid = pthread_self(); |
277 | |
|
278 | 0 | #if defined(__linux__) |
279 | | /* |
280 | | * Under Linux does not exists the difference between process and |
281 | | * threads, everything is a thread in the kernel task struct, and each |
282 | | * one has it's own numerical identificator: PID . |
283 | | * |
284 | | * Here we want to know what's the PID associated to this running |
285 | | * task (which is different from parent Monkey PID), it can be |
286 | | * retrieved with gettid() but Glibc does not export to userspace |
287 | | * the syscall, we need to call it directly through syscall(2). |
288 | | */ |
289 | 0 | worker->pid = syscall(__NR_gettid); |
290 | | #elif defined(__APPLE__) |
291 | | uint64_t tid; |
292 | | pthread_threadid_np(NULL, &tid); |
293 | | worker->pid = tid; |
294 | | #else |
295 | | worker->pid = 0xdeadbeef; |
296 | | #endif |
297 | | |
298 | | /* Initialize lists */ |
299 | 0 | mk_list_init(&worker->timeout_queue); |
300 | 0 | worker->request_handler = NULL; |
301 | |
|
302 | 0 | return worker->idx; |
303 | 0 | } |
304 | | |
305 | | static void mk_signal_thread_sigpipe_safe() |
306 | 0 | { |
307 | 0 | #ifndef _WIN32 |
308 | 0 | sigset_t old; |
309 | 0 | sigset_t set; |
310 | |
|
311 | 0 | sigemptyset(&set); |
312 | 0 | sigaddset(&set, SIGPIPE); |
313 | 0 | pthread_sigmask(SIG_BLOCK, &set, &old); |
314 | 0 | #endif |
315 | 0 | } |
316 | | |
317 | | /* created thread, all these calls are in the thread context */ |
318 | | void *mk_sched_launch_worker_loop(void *data) |
319 | 0 | { |
320 | 0 | int ret; |
321 | 0 | int wid; |
322 | 0 | unsigned long len; |
323 | 0 | char *thread_name = 0; |
324 | 0 | struct mk_list *head; |
325 | 0 | struct mk_sched_worker_cb *wcb; |
326 | 0 | struct mk_sched_worker *sched = NULL; |
327 | 0 | struct mk_sched_notif *notif = NULL; |
328 | 0 | struct mk_sched_thread_conf *thinfo = data; |
329 | 0 | struct mk_sched_ctx *ctx; |
330 | 0 | struct mk_server *server; |
331 | |
|
332 | 0 | server = thinfo->server; |
333 | 0 | ctx = server->sched_ctx; |
334 | | |
335 | | /* Avoid SIGPIPE signals on this thread */ |
336 | 0 | mk_signal_thread_sigpipe_safe(); |
337 | | |
338 | | /* Init specific thread cache */ |
339 | 0 | mk_sched_thread_lists_init(); |
340 | 0 | mk_cache_worker_init(); |
341 | | |
342 | | /* Virtual hosts: initialize per thread-vhost data */ |
343 | 0 | mk_vhost_fdt_worker_init(server); |
344 | | |
345 | | /* Register working thread */ |
346 | 0 | wid = mk_sched_register_thread(server); |
347 | 0 | sched = &ctx->workers[wid]; |
348 | 0 | sched->loop = mk_event_loop_create(MK_EVENT_QUEUE_SIZE); |
349 | 0 | if (!sched->loop) { |
350 | 0 | mk_err("Error creating Scheduler loop"); |
351 | 0 | exit(EXIT_FAILURE); |
352 | 0 | } |
353 | | |
354 | | |
355 | 0 | sched->mem_pagesize = mk_utils_get_system_page_size(); |
356 | | |
357 | | /* |
358 | | * Create the notification instance and link it to the worker |
359 | | * thread-scope list. |
360 | | */ |
361 | 0 | notif = mk_mem_alloc_z(sizeof(struct mk_sched_notif)); |
362 | 0 | MK_TLS_SET(mk_tls_sched_worker_notif, notif); |
363 | | |
364 | | /* Register the scheduler channel to signal active workers */ |
365 | 0 | ret = mk_event_channel_create(sched->loop, |
366 | 0 | &sched->signal_channel_r, |
367 | 0 | &sched->signal_channel_w, |
368 | 0 | notif); |
369 | 0 | if (ret < 0) { |
370 | 0 | exit(EXIT_FAILURE); |
371 | 0 | } |
372 | | |
373 | 0 | mk_list_init(&sched->event_free_queue); |
374 | 0 | mk_list_init(&sched->threads); |
375 | 0 | mk_list_init(&sched->threads_purge); |
376 | | |
377 | | /* |
378 | | * ULONG_MAX BUG test only |
379 | | * ======================= |
380 | | * to test the workaround we can use the following value: |
381 | | * |
382 | | * thinfo->closed_connections = 1000; |
383 | | */ |
384 | | |
385 | | //thinfo->ctx = thconf->ctx; |
386 | | |
387 | | /* Rename worker */ |
388 | 0 | mk_string_build(&thread_name, &len, "monkey: wrk/%i", sched->idx); |
389 | 0 | mk_utils_worker_rename(thread_name); |
390 | 0 | mk_mem_free(thread_name); |
391 | | |
392 | | /* Export known scheduler node to context thread */ |
393 | 0 | MK_TLS_SET(mk_tls_sched_worker_node, sched); |
394 | 0 | mk_tls_thread_init(server); |
395 | 0 | mk_plugin_core_thread(server); |
396 | |
|
397 | 0 | if (server->scheduler_mode == MK_SCHEDULER_REUSEPORT) { |
398 | 0 | sched->listeners = mk_server_listen_init(server); |
399 | 0 | if (!sched->listeners) { |
400 | 0 | exit(EXIT_FAILURE); |
401 | 0 | } |
402 | 0 | } |
403 | | |
404 | | /* Unlock the conditional initializator */ |
405 | 0 | pthread_mutex_lock(&server->pth_mutex); |
406 | 0 | server->pth_init = MK_TRUE; |
407 | 0 | pthread_cond_signal(&server->pth_cond); |
408 | 0 | pthread_mutex_unlock(&server->pth_mutex); |
409 | | |
410 | | /* Invoke custom worker-callbacks defined by the scheduler (lib) */ |
411 | 0 | mk_list_foreach(head, &server->sched_worker_callbacks) { |
412 | 0 | wcb = mk_list_entry(head, struct mk_sched_worker_cb, _head); |
413 | 0 | wcb->cb_func(wcb->data); |
414 | 0 | } |
415 | |
|
416 | 0 | mk_mem_free(thinfo); |
417 | | |
418 | | /* init server thread loop */ |
419 | 0 | mk_server_worker_loop(server); |
420 | |
|
421 | 0 | return 0; |
422 | 0 | } |
423 | | |
424 | | /* Create thread which will be listening for incomings requests */ |
425 | | int mk_sched_launch_thread(struct mk_server *server, pthread_t *tout) |
426 | 0 | { |
427 | 0 | pthread_t tid; |
428 | 0 | pthread_attr_t attr; |
429 | 0 | struct mk_sched_thread_conf *thconf; |
430 | |
|
431 | 0 | server->pth_init = MK_FALSE; |
432 | | |
433 | | /* |
434 | | * This lock is used for the 'pth_cond' conditional. Once the worker |
435 | | * thread is ready it will signal the condition. |
436 | | */ |
437 | 0 | pthread_mutex_lock(&server->pth_mutex); |
438 | | |
439 | | /* Thread data */ |
440 | 0 | thconf = mk_mem_alloc_z(sizeof(struct mk_sched_thread_conf)); |
441 | 0 | thconf->server = server; |
442 | |
|
443 | 0 | pthread_attr_init(&attr); |
444 | 0 | pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_JOINABLE); |
445 | 0 | if (pthread_create(&tid, &attr, mk_sched_launch_worker_loop, |
446 | 0 | (void *) thconf) != 0) { |
447 | 0 | mk_libc_error("pthread_create"); |
448 | 0 | pthread_mutex_unlock(&server->pth_mutex); |
449 | 0 | return -1; |
450 | 0 | } |
451 | | |
452 | 0 | *tout = tid; |
453 | | |
454 | | /* Block until the child thread is ready */ |
455 | 0 | while (!server->pth_init) { |
456 | 0 | pthread_cond_wait(&server->pth_cond, &server->pth_mutex); |
457 | 0 | } |
458 | |
|
459 | 0 | pthread_mutex_unlock(&server->pth_mutex); |
460 | |
|
461 | 0 | return 0; |
462 | 0 | } |
463 | | |
464 | | /* |
465 | | * The scheduler nodes are an array of struct mk_sched_worker type, |
466 | | * each worker thread belongs to a scheduler node, on this function we |
467 | | * allocate a scheduler node per number of workers defined. |
468 | | */ |
469 | | int mk_sched_init(struct mk_server *server) |
470 | 0 | { |
471 | 0 | int size; |
472 | 0 | struct mk_sched_ctx *ctx; |
473 | |
|
474 | 0 | ctx = mk_mem_alloc_z(sizeof(struct mk_sched_ctx)); |
475 | 0 | if (!ctx) { |
476 | 0 | mk_libc_error("malloc"); |
477 | 0 | return -1; |
478 | 0 | } |
479 | | |
480 | 0 | size = (sizeof(struct mk_sched_worker) * server->workers); |
481 | 0 | ctx->workers = mk_mem_alloc(size); |
482 | 0 | if (!ctx->workers) { |
483 | 0 | mk_libc_error("malloc"); |
484 | 0 | mk_mem_free(ctx); |
485 | 0 | return -1; |
486 | 0 | } |
487 | 0 | memset(ctx->workers, '\0', size); |
488 | | |
489 | | /* Initialize helpers */ |
490 | 0 | pthread_mutex_init(&server->pth_mutex, NULL); |
491 | 0 | pthread_cond_init(&server->pth_cond, NULL); |
492 | 0 | server->pth_init = MK_FALSE; |
493 | | |
494 | | /* Map context into server context */ |
495 | 0 | server->sched_ctx = ctx; |
496 | | |
497 | | /* The mk_thread_prepare call was replaced by mk_http_thread_initialize_tls |
498 | | * which is called earlier. |
499 | | */ |
500 | |
|
501 | 0 | return 0; |
502 | 0 | } |
503 | | |
504 | | int mk_sched_exit(struct mk_server *server) |
505 | 0 | { |
506 | 0 | struct mk_sched_ctx *ctx; |
507 | |
|
508 | 0 | ctx = server->sched_ctx; |
509 | 0 | mk_sched_worker_cb_free(server); |
510 | 0 | mk_mem_free(ctx->workers); |
511 | 0 | mk_mem_free(ctx); |
512 | |
|
513 | 0 | return 0; |
514 | 0 | } |
515 | | |
516 | | void mk_sched_set_request_list(struct rb_root *list) |
517 | 0 | { |
518 | 0 | MK_TLS_SET(mk_tls_sched_cs, list); |
519 | 0 | } |
520 | | |
521 | | int mk_sched_remove_client(struct mk_sched_conn *conn, |
522 | | struct mk_sched_worker *sched, |
523 | | struct mk_server *server) |
524 | 0 | { |
525 | 0 | struct mk_event *event; |
526 | | |
527 | | /* |
528 | | * Close socket and change status: we must invoke mk_epoll_del() |
529 | | * because when the socket is closed is cleaned from the queue by |
530 | | * the Kernel at its leisure, and we may get false events if we rely |
531 | | * on that. |
532 | | */ |
533 | 0 | event = &conn->event; |
534 | 0 | MK_TRACE("[FD %i] Scheduler remove", event->fd); |
535 | |
|
536 | 0 | mk_event_del(sched->loop, event); |
537 | | |
538 | | /* Invoke plugins in stage 50 */ |
539 | 0 | mk_plugin_stage_run_50(event->fd, server); |
540 | |
|
541 | 0 | sched->closed_connections++; |
542 | | |
543 | | /* Unlink from the red-black tree */ |
544 | | //rb_erase(&conn->_rb_head, &sched->rb_queue); |
545 | 0 | mk_sched_conn_timeout_del(conn); |
546 | | |
547 | | /* Close at network layer level */ |
548 | 0 | conn->net->close(conn->net->plugin, event->fd); |
549 | | |
550 | | /* Release and return */ |
551 | 0 | mk_channel_clean(&conn->channel); |
552 | 0 | mk_sched_event_free(&conn->event); |
553 | 0 | conn->status = MK_SCHED_CONN_CLOSED; |
554 | |
|
555 | 0 | MK_LT_SCHED(remote_fd, "DELETE_CLIENT"); |
556 | 0 | return 0; |
557 | 0 | } |
558 | | |
559 | | /* FIXME: nobody is using this function, check back later */ |
560 | | struct mk_sched_conn *mk_sched_get_connection(struct mk_sched_worker *sched, |
561 | | int remote_fd) |
562 | 0 | { |
563 | 0 | (void) sched; |
564 | 0 | (void) remote_fd; |
565 | 0 | return NULL; |
566 | 0 | } |
567 | | |
568 | | /* |
569 | | * For a given connection number, remove all resources associated: it can be |
570 | | * used on any context such as: timeout, I/O errors, request finished, |
571 | | * exceptions, etc. |
572 | | */ |
573 | | int mk_sched_drop_connection(struct mk_sched_conn *conn, |
574 | | struct mk_sched_worker *sched, |
575 | | struct mk_server *server) |
576 | 0 | { |
577 | 0 | mk_sched_threads_destroy_conn(sched, conn); |
578 | 0 | return mk_sched_remove_client(conn, sched, server); |
579 | 0 | } |
580 | | |
581 | | int mk_sched_check_timeouts(struct mk_sched_worker *sched, |
582 | | struct mk_server *server) |
583 | 0 | { |
584 | 0 | int client_timeout; |
585 | 0 | struct mk_sched_conn *conn; |
586 | 0 | struct mk_list *head; |
587 | 0 | struct mk_list *temp; |
588 | | |
589 | | /* PENDING CONN TIMEOUT */ |
590 | 0 | mk_list_foreach_safe(head, temp, &sched->timeout_queue) { |
591 | 0 | conn = mk_list_entry(head, struct mk_sched_conn, timeout_head); |
592 | 0 | if (conn->event.type & MK_EVENT_IDLE) { |
593 | 0 | continue; |
594 | 0 | } |
595 | | |
596 | 0 | client_timeout = conn->arrive_time + server->timeout; |
597 | | |
598 | | /* Check timeout */ |
599 | 0 | if (client_timeout <= server->clock_context->log_current_utime) { |
600 | 0 | MK_TRACE("Scheduler, closing fd %i due TIMEOUT", |
601 | 0 | conn->event.fd); |
602 | 0 | MK_LT_SCHED(conn->event.fd, "TIMEOUT_CONN_PENDING"); |
603 | 0 | if (conn->protocol->cb_close) { |
604 | 0 | conn->protocol->cb_close(conn, sched, MK_SCHED_CONN_TIMEOUT, |
605 | 0 | server); |
606 | 0 | } |
607 | 0 | mk_sched_drop_connection(conn, sched, server); |
608 | 0 | } |
609 | 0 | } |
610 | |
|
611 | 0 | return 0; |
612 | 0 | } |
613 | | |
614 | | static int sched_thread_cleanup(struct mk_sched_worker *sched, |
615 | | struct mk_list *list) |
616 | 0 | { |
617 | 0 | int c = 0; |
618 | 0 | struct mk_list *tmp; |
619 | 0 | struct mk_list *head; |
620 | 0 | struct mk_http_thread *mth; |
621 | 0 | (void) sched; |
622 | |
|
623 | 0 | mk_list_foreach_safe(head, tmp, list) { |
624 | 0 | mth = mk_list_entry(head, struct mk_http_thread, _head); |
625 | 0 | mk_http_thread_destroy(mth); |
626 | 0 | c++; |
627 | 0 | } |
628 | |
|
629 | 0 | return c; |
630 | |
|
631 | 0 | } |
632 | | |
633 | | int mk_sched_threads_purge(struct mk_sched_worker *sched) |
634 | 0 | { |
635 | 0 | int c = 0; |
636 | |
|
637 | 0 | c = sched_thread_cleanup(sched, &sched->threads_purge); |
638 | 0 | return c; |
639 | 0 | } |
640 | | |
641 | | int mk_sched_threads_destroy_all(struct mk_sched_worker *sched) |
642 | 0 | { |
643 | 0 | int c = 0; |
644 | |
|
645 | 0 | c = sched_thread_cleanup(sched, &sched->threads_purge); |
646 | 0 | c += sched_thread_cleanup(sched, &sched->threads); |
647 | |
|
648 | 0 | return c; |
649 | 0 | } |
650 | | |
651 | | /* |
652 | | * Destroy the thread contexts associated to the particular |
653 | | * connection. |
654 | | * |
655 | | * Return the number of contexts destroyed. |
656 | | */ |
657 | | int mk_sched_threads_destroy_conn(struct mk_sched_worker *sched, |
658 | | struct mk_sched_conn *conn) |
659 | 0 | { |
660 | 0 | int c = 0; |
661 | 0 | struct mk_list *tmp; |
662 | 0 | struct mk_list *head; |
663 | 0 | struct mk_http_thread *mth; |
664 | 0 | (void) sched; |
665 | |
|
666 | 0 | mk_list_foreach_safe(head, tmp, &sched->threads) { |
667 | 0 | mth = mk_list_entry(head, struct mk_http_thread, _head); |
668 | 0 | if (mth->session->conn == conn) { |
669 | 0 | mk_http_thread_destroy(mth); |
670 | 0 | c++; |
671 | 0 | } |
672 | 0 | } |
673 | 0 | return c; |
674 | 0 | } |
675 | | |
676 | | /* |
677 | | * Scheduler events handler: lookup for event handler and invoke |
678 | | * proper callbacks. |
679 | | */ |
680 | | static int mk_sched_conn_pending_output(struct mk_channel *channel); |
681 | | |
682 | | int mk_sched_event_read(struct mk_sched_conn *conn, |
683 | | struct mk_sched_worker *sched, |
684 | | struct mk_server *server) |
685 | 0 | { |
686 | 0 | int ret = 0; |
687 | |
|
688 | | #ifdef MK_HAVE_TRACE |
689 | | MK_TRACE("[FD %i] Connection Handler / read", conn->event.fd); |
690 | | #endif |
691 | | |
692 | | /* |
693 | | * When the event loop notify that there is some readable information |
694 | | * from the socket, we need to invoke the protocol handler associated |
695 | | * to this connection and also pass as a reference the 'read()' function |
696 | | * that handle 'read I/O' operations, e.g: |
697 | | * |
698 | | * - plain sockets through liana will use just read(2) |
699 | | * - ssl though mbedtls should use mk_mbedtls_read(..) |
700 | | */ |
701 | 0 | ret = conn->protocol->cb_read(conn, sched, server); |
702 | 0 | if (ret == -1) { |
703 | 0 | if (errno == EAGAIN) { |
704 | 0 | MK_TRACE("[FD %i] EAGAIN: need to read more data", conn->event.fd); |
705 | 0 | mk_event_add(sched->loop, conn->event.fd, |
706 | 0 | MK_EVENT_CONNECTION, |
707 | 0 | mk_net_transport_event_interest(conn->net, |
708 | 0 | conn->event.fd, |
709 | 0 | MK_EVENT_READ), |
710 | 0 | conn); |
711 | 0 | return 1; |
712 | 0 | } |
713 | 0 | return -1; |
714 | 0 | } |
715 | | |
716 | 0 | if (mk_sched_conn_pending_output(&conn->channel) == MK_TRUE) { |
717 | 0 | return mk_sched_event_write(conn, sched, server); |
718 | 0 | } |
719 | | |
720 | 0 | return ret; |
721 | 0 | } |
722 | | |
723 | | static int mk_sched_conn_pending_output(struct mk_channel *channel) |
724 | 0 | { |
725 | 0 | struct mk_list *head; |
726 | 0 | struct mk_stream *stream; |
727 | 0 | struct mk_stream_input *input; |
728 | |
|
729 | 0 | if (mk_channel_is_empty(channel) == 0) { |
730 | 0 | return MK_FALSE; |
731 | 0 | } |
732 | | |
733 | | /* |
734 | | * The request parser installs placeholder input streams with zero bytes |
735 | | * before a response exists. Only treat the channel as writable work when |
736 | | * at least one input actually has bytes pending. |
737 | | */ |
738 | 0 | mk_list_foreach(head, &channel->streams) { |
739 | 0 | stream = mk_list_entry(head, struct mk_stream, _head); |
740 | 0 | if (mk_list_is_empty(&stream->inputs) == 0) { |
741 | 0 | continue; |
742 | 0 | } |
743 | | |
744 | 0 | input = mk_list_entry_first(&stream->inputs, struct mk_stream_input, _head); |
745 | 0 | if (input->bytes_total > 0) { |
746 | 0 | return MK_TRUE; |
747 | 0 | } |
748 | 0 | } |
749 | | |
750 | 0 | return MK_FALSE; |
751 | 0 | } |
752 | | |
753 | | int mk_sched_event_write(struct mk_sched_conn *conn, |
754 | | struct mk_sched_worker *sched, |
755 | | struct mk_server *server) |
756 | 0 | { |
757 | 0 | int ret = -1; |
758 | 0 | size_t count; |
759 | 0 | struct mk_event *event; |
760 | |
|
761 | 0 | MK_TRACE("[FD %i] Connection Handler / write", conn->event.fd); |
762 | |
|
763 | 0 | ret = mk_channel_write(&conn->channel, &count); |
764 | 0 | if (ret == MK_CHANNEL_FLUSH || ret == MK_CHANNEL_BUSY) { |
765 | 0 | if ((conn->event.mask & MK_EVENT_WRITE) == 0) { |
766 | 0 | mk_event_add(sched->loop, conn->event.fd, |
767 | 0 | MK_EVENT_CONNECTION, |
768 | 0 | mk_net_transport_event_interest(conn->net, |
769 | 0 | conn->event.fd, |
770 | 0 | MK_EVENT_WRITE), |
771 | 0 | conn); |
772 | 0 | } |
773 | 0 | return 0; |
774 | 0 | } |
775 | 0 | else if (ret == MK_CHANNEL_DONE || ret == MK_CHANNEL_EMPTY) { |
776 | 0 | if (conn->protocol->cb_done) { |
777 | 0 | ret = conn->protocol->cb_done(conn, sched, server); |
778 | 0 | } |
779 | 0 | if (ret == -1) { |
780 | 0 | return -1; |
781 | 0 | } |
782 | 0 | else if (ret == 0) { |
783 | 0 | event = &conn->event; |
784 | 0 | mk_event_add(sched->loop, event->fd, |
785 | 0 | MK_EVENT_CONNECTION, |
786 | 0 | mk_net_transport_event_interest(conn->net, |
787 | 0 | event->fd, |
788 | 0 | MK_EVENT_READ), |
789 | 0 | conn); |
790 | 0 | } |
791 | 0 | return 0; |
792 | 0 | } |
793 | 0 | else if (ret & MK_CHANNEL_ERROR) { |
794 | 0 | return -1; |
795 | 0 | } |
796 | | |
797 | | /* avoid to make gcc cry :_( */ |
798 | 0 | return -1; |
799 | 0 | } |
800 | | |
801 | | int mk_sched_event_close(struct mk_sched_conn *conn, |
802 | | struct mk_sched_worker *sched, |
803 | | int type, struct mk_server *server) |
804 | 0 | { |
805 | 0 | MK_TRACE("[FD %i] Connection Handler, closed", conn->event.fd); |
806 | 0 | mk_event_del(sched->loop, &conn->event); |
807 | |
|
808 | 0 | if (type != MK_EP_SOCKET_DONE && conn->protocol->cb_close) { |
809 | 0 | conn->protocol->cb_close(conn, sched, type, server); |
810 | 0 | } |
811 | | /* |
812 | | * Remove the socket from the scheduler and make sure |
813 | | * to disable all notifications. |
814 | | */ |
815 | 0 | mk_sched_drop_connection(conn, sched, server); |
816 | 0 | return 0; |
817 | 0 | } |
818 | | |
819 | | void mk_sched_event_free(struct mk_event *event) |
820 | 0 | { |
821 | 0 | struct mk_sched_worker *sched = mk_sched_get_thread_conf(); |
822 | |
|
823 | 0 | if ((event->type & MK_EVENT_IDLE) != 0) { |
824 | 0 | return; |
825 | 0 | } |
826 | | |
827 | 0 | event->type |= MK_EVENT_IDLE; |
828 | 0 | mk_list_add(&event->_head, &sched->event_free_queue); |
829 | 0 | } |
830 | | |
831 | | /* Register a new callback function to invoke when a worker is created */ |
832 | | int mk_sched_worker_cb_add(struct mk_server *server, |
833 | | void (*cb_func) (void *), |
834 | | void *data) |
835 | 0 | { |
836 | 0 | struct mk_sched_worker_cb *wcb; |
837 | |
|
838 | 0 | wcb = mk_mem_alloc(sizeof(struct mk_sched_worker_cb)); |
839 | 0 | if (!wcb) { |
840 | 0 | return -1; |
841 | 0 | } |
842 | | |
843 | 0 | wcb->cb_func = cb_func; |
844 | 0 | wcb->data = data; |
845 | 0 | mk_list_add(&wcb->_head, &server->sched_worker_callbacks); |
846 | 0 | return 0; |
847 | 0 | } |
848 | | |
849 | | void mk_sched_worker_cb_free(struct mk_server *server) |
850 | 0 | { |
851 | 0 | struct mk_list *tmp; |
852 | 0 | struct mk_list *head; |
853 | 0 | struct mk_sched_worker_cb *wcb; |
854 | |
|
855 | 0 | mk_list_foreach_safe(head, tmp, &server->sched_worker_callbacks) { |
856 | 0 | wcb = mk_list_entry(head, struct mk_sched_worker_cb, _head); |
857 | 0 | mk_list_del(&wcb->_head); |
858 | 0 | mk_mem_free(wcb); |
859 | 0 | } |
860 | 0 | } |
861 | | |
862 | | int mk_sched_send_signal(struct mk_sched_worker *worker, uint64_t val) |
863 | 0 | { |
864 | 0 | ssize_t n; |
865 | | |
866 | | /* When using libevent _mk_event_channel_create creates a unix socket |
867 | | * instead of a pipe and windows doesn't us calling read / write on a |
868 | | * socket instead of recv / send |
869 | | */ |
870 | |
|
871 | | #ifdef _WIN32 |
872 | | n = send(worker->signal_channel_w, &val, sizeof(uint64_t), 0); |
873 | | #else |
874 | 0 | n = write(worker->signal_channel_w, &val, sizeof(uint64_t)); |
875 | 0 | #endif |
876 | |
|
877 | 0 | if (n < 0) { |
878 | 0 | mk_libc_error("write"); |
879 | |
|
880 | 0 | return 0; |
881 | 0 | } |
882 | | |
883 | 0 | return 1; |
884 | 0 | } |
885 | | |
886 | | int mk_sched_broadcast_signal(struct mk_server *server, uint64_t val) |
887 | 0 | { |
888 | 0 | int i; |
889 | 0 | int count = 0; |
890 | 0 | struct mk_sched_ctx *ctx; |
891 | 0 | struct mk_sched_worker *worker; |
892 | |
|
893 | 0 | ctx = server->sched_ctx; |
894 | 0 | for (i = 0; i < server->workers; i++) { |
895 | 0 | worker = &ctx->workers[i]; |
896 | |
|
897 | 0 | count += mk_sched_send_signal(worker, val); |
898 | 0 | } |
899 | |
|
900 | 0 | return count; |
901 | 0 | } |
902 | | |
903 | | /* |
904 | | * Wait for all workers to finish: this function assumes that previously a |
905 | | * MK_SCHED_SIGNAL_FREE_ALL was sent to the worker channels. |
906 | | */ |
907 | | int mk_sched_workers_join(struct mk_server *server) |
908 | 0 | { |
909 | 0 | int i; |
910 | 0 | int count = 0; |
911 | 0 | struct mk_sched_ctx *ctx; |
912 | 0 | struct mk_sched_worker *worker; |
913 | |
|
914 | 0 | ctx = server->sched_ctx; |
915 | 0 | for (i = 0; i < server->workers; i++) { |
916 | 0 | worker = &ctx->workers[i]; |
917 | 0 | pthread_join(worker->tid, NULL); |
918 | 0 | count++; |
919 | 0 | } |
920 | |
|
921 | 0 | return count; |
922 | 0 | } |