Coverage Report

Created: 2026-09-28 07:36

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}