vinyl-cache/bin/vinyld/cache/cache_wrk.c
0
/*-
1
 * Copyright (c) 2006 Verdens Gang AS
2
 * Copyright (c) 2006-2011 Varnish Software AS
3
 * All rights reserved.
4
 *
5
 * Author: Poul-Henning Kamp <phk@phk.freebsd.dk>
6
 *
7
 * SPDX-License-Identifier: BSD-2-Clause
8
 *
9
 * Redistribution and use in source and binary forms, with or without
10
 * modification, are permitted provided that the following conditions
11
 * are met:
12
 * 1. Redistributions of source code must retain the above copyright
13
 *    notice, this list of conditions and the following disclaimer.
14
 * 2. Redistributions in binary form must reproduce the above copyright
15
 *    notice, this list of conditions and the following disclaimer in the
16
 *    documentation and/or other materials provided with the distribution.
17
 *
18
 * THIS SOFTWARE IS PROVIDED BY THE AUTHOR AND CONTRIBUTORS ``AS IS'' AND
19
 * ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
20
 * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE
21
 * ARE DISCLAIMED.  IN NO EVENT SHALL AUTHOR OR CONTRIBUTORS BE LIABLE
22
 * FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
23
 * DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS
24
 * OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION)
25
 * HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT
26
 * LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY
27
 * OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF
28
 * SUCH DAMAGE.
29
 *
30
 * Worker thread stuff unrelated to the worker thread pools.
31
 *
32
 * --
33
 * signaling_note:
34
 *
35
 * note on worker wakeup signaling through the wrk condition variable (cv)
36
 *
37
 * In the general case, a cv needs to be signaled while holding the
38
 * corresponding mutex, otherwise the signal may be posted before the waiting
39
 * thread could register itself on the cv and, consequently, the signal may be
40
 * missed.
41
 *
42
 * In our case, any worker thread which we wake up comes from the idle queue,
43
 * where it put itself under the mutex, releasing that mutex implicitly via
44
 * Lck_CondWaitUntil() (which calls some variant of pthread_cond_wait). So we avoid
45
 * additional mutex contention knowing that any worker thread on the idle queue
46
 * is blocking on the cv.
47
 *
48
 * Except -- when it isn't, because it woke up for releasing its VCL
49
 * Reference. To account for this case, we check if the task function has been
50
 * set in the meantime, which in turn requires all of the task preparation to be
51
 * done holding the pool mutex. (see also #2719)
52
 */
53
54
#include "config.h"
55
56
#include <stdlib.h>
57
#include <sched.h>
58
59
#include "cache_int.h"
60
#include "cache_pool.h"
61
62
#include "vcli_serve.h"
63
#include "vtim.h"
64
65
#include "hash/hash_slinger.h"
66
67
static void Pool_Work_Thread(struct pool *pp, struct worker *wrk);
68
69
static uintmax_t reqpoolfail;
70
71
/*--------------------------------------------------------------------
72
 * Return all cached resources
73
 */
74
75
void
76 473500
WRK_Cleanup(const struct worker *wrk)
77
{
78
79 473500
        CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC);
80 473500
        CHECK_OBJ_NOTNULL(wrk->wpriv, WORKER_PRIV_MAGIC);
81 473500
        HSH_Cleanup(wrk);
82 473500
        if (wrk->wpriv->vcl != NULL)
83 22778
                VCL_Rel(&wrk->wpriv->vcl);
84 473500
}
85
86
/*--------------------------------------------------------------------
87
 * Create and start a back-ground thread which as its own worker and
88
 * session data structures;
89
 */
90
91
struct bgthread {
92
        unsigned        magic;
93
#define BGTHREAD_MAGIC  0x23b5152b
94
        const char      *name;
95
        bgthread_t      *func;
96
        void            *priv;
97
};
98
99
static void *
100 84210
wrk_bgthread(void *arg)
101
{
102
        struct bgthread *bt;
103
        struct worker wrk;
104
        struct worker_priv wpriv[1];
105
        struct VSC_main_wrk ds;
106
        void *r;
107
108 84210
        CAST_OBJ_NOTNULL(bt, arg, BGTHREAD_MAGIC);
109 84210
        THR_SetName(bt->name);
110 84210
        THR_Init();
111 84210
        INIT_OBJ(&wrk, WORKER_MAGIC);
112 84210
        INIT_OBJ(wpriv, WORKER_PRIV_MAGIC);
113 84210
        wrk.wpriv = wpriv;
114
        // bgthreads do not have a vpi member
115 84210
        memset(&ds, 0, sizeof ds);
116 84210
        wrk.stats = &ds;
117
118 84210
        r = bt->func(&wrk, bt->priv);
119 84210
        WRK_Cleanup(&wrk);
120 84210
        Pool_Sumstat(&wrk);
121 84210
        return (r);
122
}
123
124
void
125 84210
WRK_BgThread(pthread_t *thr, const char *name, bgthread_t *func, void *priv)
126
{
127
        struct bgthread *bt;
128
129 84210
        ALLOC_OBJ(bt, BGTHREAD_MAGIC);
130 84210
        AN(bt);
131
132 84210
        bt->name = name;
133 84210
        bt->func = func;
134 84210
        bt->priv = priv;
135 84210
        PTOK(pthread_create(thr, NULL, wrk_bgthread, bt));
136 84210
}
137
138
/*--------------------------------------------------------------------*/
139
140
static void
141 410860
WRK_Thread(struct pool *qp, size_t stacksize, unsigned thread_workspace)
142
{
143
        // child_signal_handler stack overflow check uses struct worker addr
144
        struct worker *w, ww;
145
        struct VSC_main_wrk ds;
146 410860
        unsigned char ws[thread_workspace];
147
        struct worker_priv wpriv[1];
148 410860
        unsigned char vpi[vpi_wrk_len];
149
150 410860
        AN(qp);
151 410860
        AN(stacksize);
152 410860
        AN(thread_workspace);
153
154 410860
        THR_SetName("cache-worker");
155 410860
        w = &ww;
156 410860
        INIT_OBJ(w, WORKER_MAGIC);
157 410860
        INIT_OBJ(wpriv, WORKER_PRIV_MAGIC);
158 410860
        w->wpriv = wpriv;
159 410860
        w->lastused = NAN;
160 410860
        memset(&ds, 0, sizeof ds);
161 410860
        w->stats = &ds;
162 410860
        THR_SetWorker(w);
163 410860
        PTOK(pthread_cond_init(&w->cond, NULL));
164
165 410860
        WS_Init(w->aws, "wrk", ws, thread_workspace);
166 410860
        VPI_wrk_init(w, vpi, sizeof vpi);
167 410860
        AN(w->vpi);
168
169 410860
        VSL(SLT_WorkThread, NO_VXID, "%p start", w);
170
171 410860
        Pool_Work_Thread(qp, w);
172 410860
        AZ(w->pool);
173
174 410860
        VSL(SLT_WorkThread, NO_VXID, "%p end", w);
175 410860
        PTOK(pthread_cond_destroy(&w->cond));
176 410860
        WRK_Cleanup(w);
177 410860
        Pool_Sumstat(w);
178 410860
}
179
180
/*--------------------------------------------------------------------
181
 * Summing of stats into pool counters
182
 */
183
184
static unsigned
185 740215
wrk_addstat(const struct worker *wrk, const struct pool_task *tp, unsigned locked)
186
{
187
        struct pool *pp;
188
189 740215
        CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC);
190 740215
        pp = wrk->pool;
191 740215
        CHECK_OBJ_NOTNULL(pp, POOL_MAGIC);
192 740215
        if (locked)
193 740221
                Lck_AssertHeld(&pp->mtx);
194
195 740231
        if ((tp == NULL && wrk->stats->summs > 0) ||
196 129
            (wrk->stats->summs >= cache_param->wthread_stats_rate)) {
197 1321830
                if (!locked)
198 0
                        Lck_Lock(&pp->mtx);
199
200 158374
                pp->a_stat->summs++;
201 158374
                VSC_main_Summ_wrk_wrk(pp->a_stat, wrk->stats);
202 158374
                memset(wrk->stats, 0, sizeof *wrk->stats);
203
204 158374
                if (!locked)
205 0
                        Lck_Unlock(&pp->mtx);
206 158374
        }
207
208 739973
        return (tp != NULL);
209
}
210
211
void
212 0
WRK_AddStat(const struct worker *wrk)
213
{
214
215 0
        (void)wrk_addstat(wrk, wrk->task, 0);
216 0
        wrk->stats->summs++;
217 0
}
218
219
/*--------------------------------------------------------------------
220
 * Pool reserve calculation
221
 */
222
223
static unsigned
224 899595
pool_reserve(void)
225
{
226
        unsigned lim;
227
228 899595
        if (cache_param->wthread_reserve == 0) {
229 898646
                lim = cache_param->wthread_min / 20 + 1;
230 898646
        } else {
231 949
                lim = cache_param->wthread_min * 950 / 1000;
232 949
                if (cache_param->wthread_reserve < lim)
233 640
                        lim = cache_param->wthread_reserve;
234
        }
235 899595
        if (lim < TASK_QUEUE_RESERVE)
236 891108
                return (TASK_QUEUE_RESERVE);
237 8487
        return (lim);
238 899595
}
239
240
/*--------------------------------------------------------------------*/
241
242
static struct worker *
243 159402
pool_getidleworker(struct pool *pp, enum task_prio prio)
244
{
245 159402
        struct pool_task *pt = NULL;
246
        struct worker *wrk;
247
248 159402
        CHECK_OBJ_NOTNULL(pp, POOL_MAGIC);
249 159402
        Lck_AssertHeld(&pp->mtx);
250 159402
        if (pp->nidle > (pool_reserve() * prio / TASK_QUEUE_RESERVE)) {
251 159059
                pt = VTAILQ_FIRST(&pp->idle_queue);
252 159059
                if (pt == NULL)
253 0
                        AZ(pp->nidle);
254 159059
        }
255
256 159402
        if (pt == NULL)
257 341
                return (NULL);
258
259 159061
        AZ(pt->func);
260 159061
        CAST_OBJ_NOTNULL(wrk, pt->priv, WORKER_MAGIC);
261
262 159061
        AN(pp->nidle);
263 159061
        VTAILQ_REMOVE(&pp->idle_queue, wrk->task, list);
264 159061
        pp->nidle--;
265
266 159061
        return (wrk);
267 159402
}
268
269
/*--------------------------------------------------------------------
270
 * Special scheduling:  If no thread can be found, the current thread
271
 * will be prepared for rescheduling instead.
272
 * The selected threads workspace is reserved and the argument put there.
273
 * Return one if another thread was scheduled, otherwise zero.
274
 */
275
276
int
277 49188
Pool_Task_Arg(struct worker *wrk, enum task_prio prio, task_func_t *func,
278
    const void *arg, size_t arg_len)
279
{
280
        struct pool *pp;
281
        struct worker *wrk2;
282
        int retval;
283
284 49188
        CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC);
285 49188
        AN(arg);
286 49188
        AN(arg_len);
287 49188
        pp = wrk->pool;
288 49188
        CHECK_OBJ_NOTNULL(pp, POOL_MAGIC);
289
290 49188
        Lck_Lock(&pp->mtx);
291 49188
        wrk2 = pool_getidleworker(pp, prio);
292 49188
        if (wrk2 != NULL)
293 49154
                retval = 1;
294
        else {
295 34
                wrk2 = wrk;
296 34
                retval = 0;
297
        }
298 49188
        AZ(wrk2->task->func);
299 49188
        assert(arg_len <= WS_ReserveSize(wrk2->aws, arg_len));
300 49188
        vmemcpy(WS_Reservation(wrk2->aws), arg, arg_len);
301 49188
        wrk2->task->func = func;
302 49188
        wrk2->task->priv = WS_Reservation(wrk2->aws);
303 49188
        Lck_Unlock(&pp->mtx);
304
        // see signaling_note at the top for explanation
305 49188
        if (retval)
306 49155
                PTOK(pthread_cond_signal(&wrk2->cond));
307 49188
        return (retval);
308
}
309
310
/*--------------------------------------------------------------------
311
 * Enter a new task to be done
312
 */
313
314
int
315 110233
Pool_Task(struct pool *pp, struct pool_task *task, enum task_prio prio)
316
{
317
        struct worker *wrk;
318 110233
        int retval = 0;
319 110233
        CHECK_OBJ_NOTNULL(pp, POOL_MAGIC);
320 110233
        AN(task);
321 110233
        AN(task->func);
322 110233
        assert(prio < TASK_QUEUE__END);
323
324 110233
        if (prio == TASK_QUEUE_REQ && reqpoolfail) {
325 21
                retval = reqpoolfail & 1;
326 21
                reqpoolfail >>= 1;
327 21
                if (retval) {
328 42
                        VSL(SLT_Debug, NO_VXID,
329
                            "Failing due to reqpoolfail (next= 0x%jx)",
330 21
                            reqpoolfail);
331 21
                        return (retval);
332
                }
333 0
        }
334
335 110212
        Lck_Lock(&pp->mtx);
336
337
        /* The common case first:  Take an idle thread, do it. */
338
339 110212
        wrk = pool_getidleworker(pp, prio);
340 110212
        if (wrk != NULL) {
341 109905
                AZ(wrk->task->func);
342 109905
                wrk->task->func = task->func;
343 109905
                wrk->task->priv = task->priv;
344 109905
                Lck_Unlock(&pp->mtx);
345
                // see signaling_note at the top for explanation
346 109905
                PTOK(pthread_cond_signal(&wrk->cond));
347 109905
                return (0);
348
        }
349
350
        /* Vital work is always queued. Only priority classes that can
351
         * fit under the reserve capacity are eligible to queuing.
352
         */
353 307
        if (prio >= TASK_QUEUE_RESERVE || pp->die) {
354 165
                retval = -1;
355 307
        } else if (!TASK_QUEUE_LIMITED(prio) ||
356 284
            pp->lqueue + pp->nthr < cache_param->wthread_max +
357 142
            cache_param->wthread_queue_limit) {
358 121
                pp->stats->sess_queued++;
359 121
                pp->lqueue++;
360 121
                VTAILQ_INSERT_TAIL(&pp->queues[prio], task, list);
361 121
                PTOK(pthread_cond_signal(&pp->herder_cond));
362 121
        } else {
363
                /* NB: This is counter-intuitive but when we drop a REQ
364
                 * task, it is an HTTP/1 request and we effectively drop
365
                 * the whole session. It is otherwise an h2 stream with
366
                 * STR priority in which case we are dropping a request.
367
                 */
368 21
                if (prio == TASK_QUEUE_REQ)
369 0
                        pp->stats->sess_dropped++;
370
                else
371 21
                        pp->stats->req_dropped++;
372 21
                retval = -1;
373
        }
374 307
        Lck_Unlock(&pp->mtx);
375 307
        return (retval);
376 110233
}
377
378
/*--------------------------------------------------------------------
379
 * Empty function used as a pointer value for the thread exit condition.
380
 */
381
382
static void v_matchproto_(task_func_t)
383 0
pool_kiss_of_death(struct worker *wrk, void *priv)
384
{
385 0
        (void)wrk;
386 0
        (void)priv;
387 0
}
388
389
390
/*--------------------------------------------------------------------
391
 * This is the work function for worker threads in the pool.
392
 */
393
394
static void
395 411180
Pool_Work_Thread(struct pool *pp, struct worker *wrk)
396
{
397
        struct pool_task *tp;
398
        struct pool_task tpx, tps;
399
        vtim_real tmo, now;
400
        unsigned i, reserve;
401
402 411180
        CHECK_OBJ_NOTNULL(pp, POOL_MAGIC);
403 411180
        wrk->pool = pp;
404 740234
        while (1) {
405 740234
                CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC);
406 740234
                tp = NULL;
407
408 740234
                WS_Rollback(wrk->aws, 0);
409 740234
                AZ(wrk->vsl);
410
411 740234
                Lck_Lock(&pp->mtx);
412 740234
                reserve = pool_reserve();
413
414 4017127
                for (i = 0; i < TASK_QUEUE_RESERVE; i++) {
415 3450225
                        if (pp->nidle < (reserve * i / TASK_QUEUE_RESERVE))
416 173211
                                break;
417 3277014
                        tp = VTAILQ_FIRST(&pp->queues[i]);
418 3277014
                        if (tp != NULL) {
419 121
                                pp->lqueue--;
420 121
                                pp->ndequeued--;
421 121
                                VTAILQ_REMOVE(&pp->queues[i], tp, list);
422 121
                                break;
423
                        }
424 3276893
                }
425
426 740234
                if (wrk_addstat(wrk, tp, 1)) {
427 121
                        wrk->stats->summs++;
428 121
                        AN(tp);
429 740234
                } else if (pp->b_stat != NULL && pp->a_stat->summs) {
430
                        /* Nothing to do, push pool stats into global pool */
431 158324
                        tps.func = pool_stat_summ;
432 158324
                        tps.priv = pp->a_stat;
433 158324
                        pp->a_stat = pp->b_stat;
434 158324
                        pp->b_stat = NULL;
435 158324
                        tp = &tps;
436 158324
                } else {
437
                        /* Nothing to do: To sleep, perchance to dream ... */
438 581789
                        if (isnan(wrk->lastused))
439 423430
                                wrk->lastused = VTIM_real();
440 581789
                        wrk->task->func = NULL;
441 581789
                        wrk->task->priv = wrk;
442 581789
                        VTAILQ_INSERT_HEAD(&pp->idle_queue, wrk->task, list);
443 581789
                        pp->nidle++;
444 581789
                        now = wrk->lastused;
445 581789
                        do {
446
                                // see signaling_note at the top for explanation
447 587244
                                if (DO_DEBUG(DBG_VCLREL) &&
448 42635
                                    pp->b_stat == NULL && pp->a_stat->summs)
449
                                        /* We've released the VCL, but
450
                                         * there are pool stats not pushed
451
                                         * to the global stats and some
452
                                         * thread is busy pushing
453
                                         * stats. Set a 1 second timeout
454
                                         * so that we'll wake up and get a
455
                                         * chance to push stats. */
456 36
                                        tmo = now + 1.;
457 587144
                                else if (wrk->wpriv->vcl == NULL)
458 537136
                                        tmo = INFINITY;
459 50008
                                else if (DO_DEBUG(DBG_VTC_MODE))
460 50008
                                        tmo = now + 1.;
461
                                else
462 0
                                        tmo = now + 60.;
463 587180
                                (void)Lck_CondWaitUntil(
464 587180
                                    &wrk->cond, &pp->mtx, tmo);
465 587180
                                if (wrk->task->func != NULL) {
466
                                        /* We have been handed a new task */
467 581782
                                        tpx = *wrk->task;
468 581782
                                        tp = &tpx;
469 581782
                                        wrk->stats->summs++;
470 587180
                                } else if (pp->b_stat != NULL &&
471 5392
                                    pp->a_stat->summs) {
472
                                        /* Woken up to release the VCL,
473
                                         * and noticing that there are
474
                                         * pool stats not pushed to the
475
                                         * global stats and no active
476
                                         * thread currently doing
477
                                         * it. Remove ourself from the
478
                                         * idle queue and take on the
479
                                         * task. */
480 7
                                        assert(pp->nidle > 0);
481 7
                                        VTAILQ_REMOVE(&pp->idle_queue,
482
                                            wrk->task, list);
483 7
                                        pp->nidle--;
484 7
                                        tps.func = pool_stat_summ;
485 7
                                        tps.priv = pp->a_stat;
486 7
                                        pp->a_stat = pp->b_stat;
487 7
                                        pp->b_stat = NULL;
488 7
                                        tp = &tps;
489 7
                                } else {
490
                                        // Presumably ETIMEDOUT but we do not
491
                                        // assert this because pthread condvars
492
                                        // are not airtight.
493 5391
                                        if (wrk->wpriv->vcl)
494 5367
                                                VCL_Rel(&wrk->wpriv->vcl);
495 5391
                                        now = VTIM_real();
496
                                }
497 587180
                        } while (tp == NULL);
498
                }
499 740234
                Lck_Unlock(&pp->mtx);
500
501 740234
                if (tp->func == pool_kiss_of_death)
502 411180
                        break;
503
504 329054
                do {
505 424546
                        memset(wrk->task, 0, sizeof wrk->task);
506 424546
                        assert(wrk->pool == pp);
507 424546
                        AN(tp->func);
508 424546
                        tp->func(wrk, tp->priv);
509 424546
                        if (DO_DEBUG(DBG_VCLREL) && wrk->wpriv->vcl != NULL)
510 374
                                VCL_Rel(&wrk->wpriv->vcl);
511 424546
                        tpx = *wrk->task;
512 424546
                        tp = &tpx;
513 424546
                } while (tp->func != NULL);
514
515 329054
                if (WS_Overflowed(wrk->aws))
516 21
                        wrk->stats->ws_thread_overflow++;
517
                /* cleanup for next task */
518 329054
                wrk->seen_methods = 0;
519
        }
520 411180
        wrk->pool = NULL;
521 411180
        Lck_Lock(&pp->mtx);
522 411180
        AN(pp->wrk_dying);
523 411180
        pp->wrk_dying--;
524 411180
        Lck_Unlock(&pp->mtx);
525 411180
}
526
527
/*--------------------------------------------------------------------
528
 * Create another worker thread.
529
 */
530
531
struct pool_info {
532
        unsigned                magic;
533
#define POOL_INFO_MAGIC         0x4e4442d3
534
        size_t                  stacksize;
535
        struct pool             *qp;
536
};
537
538
static void *
539 411548
pool_thread(void *priv)
540
{
541
        struct pool_info *pi;
542
543 411548
        CAST_OBJ_NOTNULL(pi, priv, POOL_INFO_MAGIC);
544 411548
        THR_Init();
545 411548
        WRK_Thread(pi->qp, pi->stacksize, cache_param->workspace_thread);
546 411548
        FREE_OBJ(pi);
547 411548
        return (NULL);
548
}
549
550
static void
551 423509
pool_breed(struct pool *qp)
552
{
553
        pthread_t tp;
554
        pthread_attr_t tp_attr;
555
        struct pool_info *pi;
556
557 423509
        PTOK(pthread_attr_init(&tp_attr));
558 423509
        PTOK(pthread_attr_setdetachstate(&tp_attr, PTHREAD_CREATE_DETACHED));
559
560
        /* Set the stacksize for worker threads we create */
561 423509
        if (cache_param->wthread_stacksize != UINT_MAX)
562 423509
                PTOK(pthread_attr_setstacksize(&tp_attr, cache_param->wthread_stacksize));
563
564 423509
        ALLOC_OBJ(pi, POOL_INFO_MAGIC);
565 423509
        AN(pi);
566 423509
        PTOK(pthread_attr_getstacksize(&tp_attr, &pi->stacksize));
567 423509
        pi->qp = qp;
568
569 423509
        errno = pthread_create(&tp, &tp_attr, pool_thread, pi);
570 423509
        if (errno) {
571 0
                FREE_OBJ(pi);
572 0
                VSL(SLT_Debug, NO_VXID, "Create worker thread failed %d %s",
573 0
                    errno, VAS_errtxt(errno));
574 0
                Lck_Lock(&pool_mtx);
575 0
                VSC_C_main->threads_failed++;
576 0
                Lck_Unlock(&pool_mtx);
577 0
                VTIM_sleep(cache_param->wthread_fail_delay);
578 0
        } else {
579 423509
                qp->nthr++;
580 423509
                Lck_Lock(&pool_mtx);
581 423509
                VSC_C_main->threads++;
582 423509
                VSC_C_main->threads_created++;
583 423509
                Lck_Unlock(&pool_mtx);
584 423509
                if (cache_param->wthread_add_delay > 0.0)
585 324
                        VTIM_sleep(cache_param->wthread_add_delay);
586
                else
587 423185
                        (void)sched_yield();
588
        }
589
590 423509
        PTOK(pthread_attr_destroy(&tp_attr));
591 423509
}
592
593
/*--------------------------------------------------------------------
594
 * Herd a single pool
595
 *
596
 * This thread wakes up every thread_pool_timeout seconds, whenever a pool
597
 * queues and when threads need to be destroyed
598
 *
599
 * The trick here is to not be too aggressive about creating threads.  In
600
 * pool_breed(), we sleep whenever we create a thread and a little while longer
601
 * whenever we fail to, hopefully missing a lot of cond_signals in the meantime.
602
 *
603
 * Idle threads are destroyed at a rate determined by wthread_destroy_delay
604
 *
605
 * XXX: probably need a lot more work.
606
 *
607
 */
608
609
void*
610 41480
pool_herder(void *priv)
611
{
612
        struct pool *pp;
613
        struct pool_task *pt;
614
        double t_idle;
615
        struct worker *wrk;
616
        double delay;
617
        unsigned wthread_min;
618 41480
        uintmax_t dq = (1ULL << 31);
619 41480
        vtim_mono dqt = 0;
620 41480
        int r = 0;
621
622 41480
        CAST_OBJ_NOTNULL(pp, priv, POOL_MAGIC);
623
624 41480
        THR_SetName("pool_herder");
625 41480
        THR_Init();
626
627 1012678
        while (!pp->die || pp->nthr > 0) {
628
                /*
629
                 * If the worker pool is configured too small, we can
630
                 * end up deadlocking it (see #2418 for details).
631
                 *
632
                 * Recovering from this would require a lot of complicated
633
                 * code, and fundamentally, either people configured their
634
                 * pools wrong, in which case we want them to notice, or
635
                 * they are under DoS, in which case recovering gracefully
636
                 * is unlikely be a major improvement.
637
                 *
638
                 * Instead we implement a watchdog and kill the worker if
639
                 * nothing has been dequeued for that long.
640
                 */
641 970692
                if (VTAILQ_EMPTY(&pp->queues[TASK_QUEUE_HIGHEST_PRIORITY])) {
642
                        /* Watchdog only applies to no movement on the
643
                         * highest priority queue (TASK_QUEUE_BO) */
644 970630
                        dq = pp->ndequeued + 1;
645 970692
                } else if (dq != pp->ndequeued) {
646 62
                        dq = pp->ndequeued;
647 62
                        dqt = VTIM_mono();
648 62
                } else if (VTIM_mono() - dqt > cache_param->wthread_watchdog) {
649 0
                        VSL(SLT_Error, NO_VXID,
650
                            "Pool Herder: Queue does not move ql=%u dt=%f",
651 0
                            pp->lqueue, VTIM_mono() - dqt);
652 0
                        WRONG("Worker Pool Queue does not move"
653
                              " - see thread_pool_watchdog parameter");
654 0
                }
655 970692
                wthread_min = cache_param->wthread_min;
656 970692
                if (pp->die)
657 411338
                        wthread_min = 0;
658
659
                /* Make more threads if needed and allowed */
660 970869
                if (pp->nthr < wthread_min ||
661 547435
                    (pp->lqueue > 0 && pp->nthr < cache_param->wthread_max)) {
662 423552
                        pool_breed(pp);
663 423552
                        continue;
664
                }
665
666 547730
                delay = cache_param->wthread_timeout;
667 547730
                assert(pp->nthr >= wthread_min);
668
669 547730
                if (pp->nthr > wthread_min) {
670
671 410939
                        t_idle = VTIM_real() - cache_param->wthread_timeout;
672
673 410939
                        Lck_Lock(&pp->mtx);
674 410939
                        wrk = NULL;
675 410939
                        pt = VTAILQ_LAST(&pp->idle_queue, taskhead);
676 410939
                        if (pt != NULL) {
677 410764
                                AN(pp->nidle);
678 410764
                                AZ(pt->func);
679 410764
                                CAST_OBJ_NOTNULL(wrk, pt->priv, WORKER_MAGIC);
680
681 410764
                                if (pp->die || wrk->lastused < t_idle ||
682 60
                                    pp->nthr > cache_param->wthread_max) {
683
                                        /* Give it a kiss on the cheek... */
684 410704
                                        VTAILQ_REMOVE(&pp->idle_queue,
685
                                            wrk->task, list);
686 410704
                                        pp->nidle--;
687 410704
                                        wrk->task->func = pool_kiss_of_death;
688 410704
                                        PTOK(pthread_cond_signal(&wrk->cond));
689 410704
                                        pp->nthr--;
690 410704
                                        pp->wrk_dying++;
691 410704
                                } else {
692 60
                                        delay = wrk->lastused - t_idle;
693 60
                                        wrk = NULL;
694
                                }
695 410764
                        }
696 410939
                        Lck_Unlock(&pp->mtx);
697
698 410939
                        if (wrk != NULL) {
699 410659
                                Lck_Lock(&pool_mtx);
700 410659
                                VSC_C_main->threads--;
701 410659
                                VSC_C_main->threads_destroyed++;
702 410659
                                Lck_Unlock(&pool_mtx);
703 410659
                                delay = cache_param->wthread_destroy_delay;
704 410659
                        } else
705 196
                                delay = vmax(delay,
706
                                    cache_param->wthread_destroy_delay);
707 410855
                }
708
709 547646
                if (pp->die) {
710 411653
                        if (delay < 2)
711 411528
                                delay = .01;
712
                        else
713 125
                                delay = 1;
714 411653
                        VTIM_sleep(delay);
715 411653
                        continue;
716
                }
717 135993
                Lck_Lock(&pp->mtx);
718 135993
                if (pp->lqueue == 0) {
719 135930
                        if (DO_DEBUG(DBG_VTC_MODE))
720 135856
                                delay = 0.5;
721 135930
                        r = Lck_CondWaitTimeout(
722 135930
                            &pp->herder_cond, &pp->mtx, delay);
723 135993
                } else if (pp->nthr >= cache_param->wthread_max) {
724
                        /* XXX: unsafe counters */
725 63
                        if (r != ETIMEDOUT)
726 21
                                VSC_C_main->threads_limited++;
727 63
                        r = Lck_CondWaitTimeout(
728 63
                            &pp->herder_cond, &pp->mtx, 1.0);
729 63
                }
730 135993
                Lck_Unlock(&pp->mtx);
731
        }
732 41986
        return (NULL);
733
}
734
735
/*--------------------------------------------------------------------
736
 * Debugging aids
737
 */
738
739
static void v_matchproto_(cli_func_t)
740 21
debug_reqpoolfail(struct cli *cli, const char * const *av, void *priv)
741
{
742 21
        uintmax_t u = 1;
743
        const char *p;
744
745 21
        (void)priv;
746 21
        (void)cli;
747 21
        reqpoolfail = 0;
748 42
        for (p = av[2]; *p != '\0'; p++) {
749 21
                if (*p == 'F' || *p == 'f')
750 21
                        reqpoolfail |= u;
751 21
                u <<= 1;
752 21
        }
753 21
}
754
755
static struct cli_proto debug_cmds[] = {
756
        { CLICMD_DEBUG_REQPOOLFAIL,             "d", debug_reqpoolfail },
757
        { NULL }
758
};
759
760
void
761 115401
WRK_Log(enum VSL_tag_e tag, const char *fmt, ...)
762
{
763
        struct worker *wrk;
764
        va_list ap;
765
766 115401
        AN(fmt);
767
768 115401
        wrk = THR_GetWorker();
769 115401
        CHECK_OBJ_ORNULL(wrk, WORKER_MAGIC);
770
771 115401
        va_start(ap, fmt);
772 115401
        if (wrk != NULL && wrk->vsl != NULL)
773 70742
                VSLbv(wrk->vsl, tag, fmt, ap);
774
        else
775 44659
                VSLv(tag, NO_VXID, fmt, ap);
776 115401
        va_end(ap);
777 115401
}
778
779
/*--------------------------------------------------------------------
780
 *
781
 */
782
783
void
784 20888
WRK_Init(void)
785
{
786 20888
        assert(cache_param->wthread_min >= TASK_QUEUE_RESERVE);
787 20888
        CLI_AddFuncs(debug_cmds);
788 20888
}