The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
worker.c
Go to the documentation of this file.
1/*
2 * This program is free software; you can redistribute it and/or modify
3 * it under the terms of the GNU General Public License as published by
4 * the Free Software Foundation; either version 2 of the License, or
5 * (at your option) any later version.
6 *
7 * This program is distributed in the hope that it will be useful,
8 * but WITHOUT ANY WARRANTY; without even the implied warranty of
9 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
10 * GNU General Public License for more details.
11 *
12 * You should have received a copy of the GNU General Public License
13 * along with this program; if not, write to the Free Software
14 * Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301, USA
15 */
16
17 /**
18 * $Id: cdd76e35a472d930710d3785bd0f0384d63aac9c $
19 *
20 * @brief Worker thread functions.
21 * @file io/worker.c
22 *
23 * The "worker" thread is the one responsible for the bulk of the
24 * work done when processing a request. Workers are spawned by the
25 * scheduler, and create a kqueue (KQ) and control-plane
26 * Atomic Queue (AQ) for control-plane communication.
27 *
28 * When a network thread discovers that it needs more workers, it
29 * asks the scheduler for a KQ/AQ combination. The network thread
30 * then creates a channel dedicated to that worker, and sends the
31 * channel to the worker in a "new channel" message. The worker
32 * receives the channel, and sends an ACK back to the network thread.
33 *
34 * The network thread then sends the worker new packets, which the
35 * worker receives and processes.
36 *
37 * When a packet is decoded, it is put into the "runnable" heap, and
38 * also into the timeout sublist. The main loop fr_worker() then
39 * pulls new requests off of this heap and runs them. The main event
40 * loop checks the head of the timeout sublist, and forcefully terminates
41 * any requests which have been running for too long.
42 *
43 * If a request is yielded, it is placed onto the yielded list in
44 * the worker "tracking" data structure.
45 *
46 * @copyright 2016 Alan DeKok (aland@freeradius.org)
47 */
48
49RCSID("$Id: cdd76e35a472d930710d3785bd0f0384d63aac9c $")
50
51#define LOG_PREFIX worker->name
52#define LOG_DST worker->log
53
54#include <freeradius-devel/io/channel.h>
55#include <freeradius-devel/io/listen.h>
56#include <freeradius-devel/io/worker.h>
57#include <freeradius-devel/unlang/base.h>
58#include <freeradius-devel/util/minmax_heap.h>
59#include <freeradius-devel/util/timer.h>
60
61#include <stdalign.h>
62
63#ifdef WITH_VERIFY_PTR
64static void worker_verify(fr_worker_t *worker);
65#define WORKER_VERIFY worker_verify(worker)
66#else
67#define WORKER_VERIFY
68#endif
69
70static _Atomic(uint64_t) request_number = 0;
71
74
75static _Thread_local fr_ring_buffer_t *fr_worker_rb;
76
77typedef struct {
78 fr_worker_t *worker; //!< the worker that owns this channel slot,
79 ///< so channel callbacks can reach back without
80 ///< threading `worker` through every layer.
81 fr_channel_t *ch;
82 fr_message_set_t *ms; //!< messages for this channel
83
84 fr_dlist_head_t dlist; //!< of requests received on this channel
86
87/**
88 * A worker which takes packets from a master, and processes them.
89 */
91 char const *name; //!< name of this worker
92 fr_worker_config_t config; //!< external configuration
93
94 unlang_interpret_t *intp; //!< Worker's local interpreter.
95
96 pthread_t thread_id; //!< my thread ID
97
98 fr_log_t const *log; //!< log destination
99 fr_log_lvl_t lvl; //!< log level
100
101 fr_atomic_queue_t *aq_control; //!< atomic queue for control messages sent to me
102
103 fr_control_t *control; //!< the control plane
104
105 fr_event_list_t *el; //!< our event list
106
107 int num_channels; //!< actual number of channels
108
109 fr_heap_t *runnable; //!< current runnable requests which we've spent time processing
110
111 fr_timer_list_t *timeout; //!< Track when requests timeout using a dlist.
112 fr_time_delta_t max_request_time; //!< maximum time a request can be processed
113
114 fr_rb_tree_t *dedup; //!< de-dup tree
115
116 fr_rb_tree_t *listeners; //!< so we can cancel requests when a listener goes away
117
118 fr_io_stats_t stats; //!< input / output stats
119 fr_time_elapsed_t cpu_time; //!< histogram of total CPU time per request
120 fr_time_elapsed_t wall_clock; //!< histogram of wall clock time per request
121
122 uint64_t num_naks; //!< number of messages which were nak'd
123 uint64_t num_active; //!< number of active requests
124
125 fr_time_delta_t predicted; //!< How long we predict a request will take to execute.
126 fr_time_tracking_t tracking; //!< how much time the worker has spent doing things.
127
128 bool was_sleeping; //!< used to suppress multiple sleep signals in a row
129 bool exiting; //!< are we exiting?
130
131 fr_worker_channel_t *channel; //!< list of channels
132
133 request_slab_list_t *slab; //!< slab allocator for request_t
134};
135
136typedef struct {
137 fr_listen_t const *listener; //!< incoming packets
138
139 fr_rb_node_t node; //!< in tree of listeners
140
141 /*
142 * To save time, we don't care about num_elements here. Which means that we don't
143 * need to cache or lookup the fr_worker_listen_t when we free a request.
144 */
145 fr_dlist_head_t dlist; //!< of requests associated with this listener.
147
148
149static fr_cmp_ret_t worker_listener_cmp(void const *one, void const *two)
150{
151 fr_worker_listen_t const *a = one, *b = two;
152
153 return CMP(a->listener, b->listener);
154}
155
156
157/*
158 * Explicitly cleanup the memory allocated to the ring buffer,
159 * just in case valgrind complains about it.
160 */
161static int _fr_worker_rb_free(void *arg)
162{
163 return talloc_free(arg);
164}
165
166/** Initialise thread local storage
167 *
168 * @return fr_ring_buffer_t for messages
169 */
171{
173
174 rb = fr_worker_rb;
175 if (rb) return rb;
176
178 if (!rb) {
179 fr_perror("Failed allocating memory for worker ring buffer");
180 return NULL;
181 }
182
184
185 return rb;
186}
187
188static inline bool is_worker_thread(fr_worker_t const *worker)
189{
190 return (pthread_equal(pthread_self(), worker->thread_id) != 0);
191}
192
194static void worker_send_reply(fr_worker_t *worker, request_t *request, bool do_not_respond, fr_time_t now);
195
196/** Callback which handles a message being received on the worker side.
197 *
198 * @param[in] ch the channel to drain
199 * @param[in] cd the message (if any) to start with
200 * @param[in] uctx the worker channel slot the message came in on
201 */
202static void worker_recv_request(fr_channel_t *ch, fr_channel_data_t *cd, void *uctx)
203{
204 fr_worker_channel_t *wc = uctx;
205 fr_worker_t *worker = wc->worker;
206
207 worker->stats.in++;
208 DEBUG3("Received request %" PRIu64 "", worker->stats.in);
209 cd->channel.ch = ch;
211}
212
214{
215 fr_async_t *async;
216
217 while ((async = fr_dlist_pop_head(&ch->dlist)) != NULL) {
219 }
220}
221
222static void worker_exit(fr_worker_t *worker)
223{
224 worker->exiting = true;
225
226 /*
227 * Don't allow the post event to run
228 * any more requests. They'll be
229 * signalled to stop before we exit.
230 *
231 * This only has an effect in single
232 * threaded mode.
233 */
234 (void)fr_event_post_delete(worker->el, fr_worker_post_event, worker);
235}
236
237/** Handle a control plane message sent to the worker via a channel
238 *
239 * @param[in] data the message
240 * @param[in] data_size size of the data
241 * @param[in] now the current time
242 * @param[in] uctx the worker
243 */
244static void worker_channel_callback(void const *data, size_t data_size, fr_time_t now, void *uctx)
245{
246 int i;
247 unsigned int num;
248 bool ok, was_sleeping;
249 fr_channel_t *ch;
252 fr_worker_t *worker = uctx;
253
254 was_sleeping = worker->was_sleeping;
255 worker->was_sleeping = false;
256
257 /*
258 * We were woken up by a signal to do something. We're
259 * not sleeping.
260 */
261 ce = fr_channel_service_message(now, &ch, data, data_size);
262 DEBUG3("Channel %s",
263 fr_table_str_by_value(channel_signals, ce, "<INVALID>"));
264 switch (ce) {
265 case FR_CHANNEL_ERROR:
266 return;
267
268 case FR_CHANNEL_EMPTY:
269 return;
270
271 case FR_CHANNEL_NOOP:
272 return;
273
275 fr_assert(0 == 1);
276 break;
277
279 fr_assert(ch != NULL);
280
281 if (!fr_channel_recv_request(ch)) {
282 worker->was_sleeping = was_sleeping;
283
284 } else while (fr_channel_recv_request(ch));
285 break;
286
287 case FR_CHANNEL_OPEN:
288 fr_assert(ch != NULL);
289
290 ok = false;
291 for (i = 0; i < worker->config.max_channels; i++) {
292 fr_assert(worker->channel[i].ch != ch);
293
294 if (worker->channel[i].ch != NULL) continue;
295
296 worker->channel[i].worker = worker;
297 worker->channel[i].ch = ch;
298 fr_dlist_init(&worker->channel[i].dlist, fr_async_t, entry);
299
300 DEBUG3("Received channel %p into array entry %d", ch, i);
301
302 ms = fr_message_set_create(worker, worker->config.message_set_size,
303 sizeof(fr_channel_data_t),
304 worker->config.ring_buffer_size, false);
305 fr_assert(ms != NULL);
306 worker->channel[i].ms = ms;
307
308 /*
309 * Hand the channel a pointer to the slot rather than to
310 * any one field of it, so a callback holding only the
311 * channel reaches the message set and the request list
312 * alike. Neither can be set in
313 * fr_worker_channel_create() because the slot has not
314 * been claimed at that point.
315 */
316 fr_channel_responder_uctx_add(ch, &worker->channel[i]);
318
319 worker->num_channels++;
320 ok = true;
321 break;
322 }
323
324 fr_cond_assert(ok);
325 break;
326
327 case FR_CHANNEL_CLOSE:
328 fr_assert(ch != NULL);
329
330 ok = false;
331
332 /*
333 * Locate the signalling channel in the list
334 * of channels.
335 */
336 for (i = 0; i < worker->config.max_channels; i++) {
337 if (!worker->channel[i].ch) continue;
338
339 if (worker->channel[i].ch != ch) continue;
340
341 worker_requests_cancel(&worker->channel[i]);
342
343 ms = worker->channel[i].ms;
344
345 fr_assert_msg(fr_dlist_num_elements(&worker->channel[i].dlist) == 0,
346 "Network added messages to channel after sending FR_CHANNEL_CLOSE");
347
348 /*
349 * Should be nothing left: the network is not supposed
350 * to enqueue anything once it has signalled the close,
351 * which is what the assert above claims. Hand back
352 * whatever we find anyway, so the messages do not
353 * strand the ring buffer they came from, and complain,
354 * because these produce no reply and so the network
355 * never decrements its outstanding count for them.
356 */
358 if (num > 0) PWARN("Discarded %u request(s) still queued at close", num);
359
361 fr_assert(ms != NULL);
363 talloc_free(ms);
364
365 worker->channel[i].ch = NULL;
366
367 fr_assert(fr_dlist_num_elements(&worker->channel[i].dlist) == 0);
368 fr_assert(worker->num_channels > 0);
369
370 worker->num_channels--;
371 ok = true;
372 break;
373 }
374
375 fr_cond_assert(ok);
376
377 /*
378 * Our last input channel closed,
379 * time to die.
380 */
381 if (worker->num_channels == 0) worker_exit(worker);
382 break;
383 }
384}
385
387{
389 request_t *request;
390
391 fr_rb_find((void **)&wl, worker->listeners, &(fr_worker_listen_t) { .listener = li });
392 if (!wl) return -1;
393
394 while ((request = fr_dlist_pop_head(&wl->dlist)) != NULL) {
395 RERROR("Cancelling request due to socket being closed");
397 }
398
399 (void) fr_rb_delete(worker->listeners, wl);
400 talloc_free(wl);
401
402 return 0;
403}
404
405
406/** A socket is going away, so clean up any requests which use this socket.
407 *
408 * @param[in] data the message
409 * @param[in] data_size size of the data
410 * @param[in] now the current time
411 * @param[in] uctx the worker
412 */
413static void worker_listen_cancel_callback(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
414{
415 fr_listen_t const *li;
416 fr_worker_t *worker = uctx;
417
418 fr_assert(data_size == sizeof(li));
419
420 memcpy(&li, data, sizeof(li));
421
422 (void) fr_worker_listen_cancel_self(worker, li);
423}
424
425/** Send a NAK to the network thread
426 *
427 * The network thread believes that a worker is running a request until that request has been NAK'd.
428 * We typically NAK requests when they've been hanging around in the worker's backlog too long,
429 * or there was an error executing the request.
430 *
431 * @param[in] worker the worker
432 * @param[in] cd the message to NAK
433 * @param[in] now when the message is NAKd
434 */
435static void worker_nak(fr_worker_t *worker, fr_channel_data_t *cd, fr_time_t now)
436{
437 size_t size;
438 fr_channel_data_t *reply;
439 fr_channel_t *ch;
442 fr_listen_t *listen;
443
444 worker->num_naks++;
445
446 /*
447 * Cache the outbound channel. We'll need it later.
448 */
449 ch = cd->channel.ch;
450 listen = cd->listen;
451
452 /*
453 * If the channel has been closed, but we haven't
454 * been informed, that is extremely bad.
455 *
456 * Try to continue working... but we'll likely
457 * leak memory or SEGV soon.
458 */
459 if (!fr_cond_assert_msg(fr_channel_active(ch), "Wanted to send NAK but channel has been closed")) {
460 fr_message_done(&cd->m);
461 return;
462 }
463
465 ms = wc->ms;
466 fr_assert(ms != NULL);
467
468 size = listen->app_io->default_reply_size;
469 if (!size) size = listen->app_io->default_message_size;
470
471 /*
472 * Allocate a default message size.
473 */
475
476 /*
477 * Encode a NAK
478 */
479 if (listen->app_io->nak) {
480 size = listen->app_io->nak(listen, cd->packet_ctx, cd->m.data,
481 cd->m.data_size, reply->m.data, reply->m.rb_size);
482 } else {
483 size = 1; /* rely on them to figure it the heck out */
484 }
485
486 (void) fr_message_and_data_commit(ms, &reply->m, size);
487
488 /*
489 * Fill in the NAK.
490 */
491 reply->m.when = now;
492 reply->reply.cpu_time = worker->tracking.running_total;
493 reply->reply.processing_time = fr_time_delta_from_msec(1); /* @todo - set to something better? */
494 reply->reply.request_time = cd->request.recv_time;
495
496 reply->listen = cd->listen;
497 reply->packet_ctx = cd->packet_ctx;
498
499 /*
500 * Mark the original message as done.
501 */
502 fr_message_done(&cd->m);
503
504 /*
505 * Send the reply, which also polls the request queue.
506 */
507 if (fr_channel_send_reply(ch, reply) < 0) {
508 DEBUG2("Failed sending reply to channel");
509 }
510
511 worker->stats.out++;
512}
513
514/** Signal the unlang interpreter that it needs to stop running the request
515 *
516 * Signalling is a synchronous operation. Whatever I/O requests the request
517 * is currently performing are immediately cancelled, and all the frames are
518 * popped off the unlang stack.
519 *
520 * Modules and unlang keywords explicitly register signal handlers to deal
521 * with their yield points being cancelled/interrupted via this function.
522 *
523 * The caller should assume the request is no longer viable after calling
524 * this function.
525 *
526 * @param[in] request request to cancel. The request may still run to completion.
527 */
528static void worker_stop_request(request_t *request)
529{
530 /*
531 * Also marks the request as done and runs
532 * the internal/external callbacs.
533 */
535}
536
537/** Enforce max_request_time
538 *
539 * Run periodically, and tries to clean up requests which were received by the network
540 * thread more than max_request_time seconds ago. In the interest of not adding a
541 * timer for every packet, the requests are given a 1 second leeway.
542 *
543 * @param[in] tl the worker's timer list.
544 * @param[in] when the current time
545 * @param[in] uctx the request_t timing out.
546 */
548{
549 request_t *request = talloc_get_type_abort(uctx, request_t);
550
551 /*
552 * Waiting too long, delete it.
553 */
554 REDEBUG("Request has reached max_request_time - signalling it to stop");
555 worker_stop_request(request);
556
557 /*
558 * This ensures the finally section can run timeout specific policies
559 */
560 request->rcode = RLM_MODULE_TIMEOUT;
561}
562
563
564/** Start time tracking for a request, and mark it as runnable.
565 *
566 */
568{
569 /*
570 * New requests are inserted into the time order heap in
571 * strict time priority. Once they are in the list, they
572 * are only removed when the request is done / free'd.
573 */
574 fr_assert(!fr_timer_armed(request->timeout));
575
576 if (unlikely(fr_timer_in(request, worker->timeout, &request->timeout, worker->config.max_request_time,
577 true, _worker_request_timeout, request) < 0)) {
578 RERROR("Failed to set request timeout timer");
579 return -1;
580 }
581
582 /*
583 * Bootstrap the async state machine with the initial
584 * state of the request.
585 */
586 RDEBUG3("Time tracking started in yielded state");
587 fr_time_tracking_start(&worker->tracking, &request->async->tracking, now);
588 fr_time_tracking_yield(&request->async->tracking, now);
589 worker->num_active++;
590
591 fr_assert(!fr_heap_entry_inserted(request->runnable));
592 (void) fr_heap_insert(&worker->runnable, request);
593
594 return 0;
595}
596
598{
599 RDEBUG3("Time tracking ended");
600 fr_time_tracking_end(&worker->predicted, &request->async->tracking, now);
601 fr_assert(worker->num_active > 0);
602 worker->num_active--;
603
604 TALLOC_FREE(request->timeout); /* Disarm the reques timer */
605}
606
607/** Send a response packet to the network side
608 *
609 * @param[in] worker This worker.
610 * @param[in] request we're sending a reply for.
611 * @param[in] send_reply whether the network side sends a reply
612 * @param[in] now The current time
613 */
614static void worker_send_reply(fr_worker_t *worker, request_t *request, bool send_reply, fr_time_t now)
615{
616 fr_channel_data_t *reply;
617 fr_channel_t *ch;
620 size_t size = 1;
621
622 REQUEST_VERIFY(request);
623
624 /*
625 * If we're sending a reply, then it's no longer runnable.
626 */
627 fr_assert(!fr_heap_entry_inserted(request->runnable));
628
629 if (send_reply) {
630 size = request->async->listen->app_io->default_reply_size;
631 if (!size) size = request->async->listen->app_io->default_message_size;
632 }
633
634 /*
635 * Allocate and send the reply.
636 */
637 ch = request->async->channel;
638 fr_assert(ch != NULL);
639
640 /*
641 * If the channel has been closed, but we haven't
642 * been informed, that is extremely bad.
643 *
644 * Try to continue working... but we'll likely
645 * leak memory or SEGV soon.
646 */
647 if (!fr_cond_assert_msg(fr_channel_active(ch), "Wanted to send reply but channel has been closed")) {
648 return;
649 }
650
652 ms = wc->ms;
653 fr_assert(ms != NULL);
654
656 fr_assert(reply != NULL);
657
658 /*
659 * Encode it, if required.
660 */
661 if (send_reply) {
662 ssize_t slen = 0;
663 fr_listen_t const *listen = request->async->listen;
664
665 if (listen->app_io->encode) {
666 slen = listen->app_io->encode(listen->app_io_instance, request,
667 reply->m.data, reply->m.rb_size);
668 } else if (listen->app->encode) {
669 slen = listen->app->encode(listen->app_instance, request,
670 reply->m.data, reply->m.rb_size);
671 }
672 if (slen < 0) {
673 RPERROR("Failed encoding request");
674 *reply->m.data = 0;
675 slen = 1;
676 }
677
678 /*
679 * Shrink the buffer to the actual packet size.
680 *
681 * This will ALWAYS return the same message as we put in.
682 */
683 fr_assert((size_t) slen <= reply->m.rb_size);
684 (void) fr_message_and_data_commit(ms, &reply->m, slen);
685 } else {
686 (void) fr_message_and_data_commit(ms, &reply->m, 0);
687 }
688
689 /*
690 * Fill in the rest of the fields in the channel message.
691 *
692 * sequence / ack will be filled in by fr_channel_send_reply()
693 */
694 reply->m.when = now;
695 reply->reply.cpu_time = worker->tracking.running_total;
696 reply->reply.processing_time = request->async->tracking.running_total;
697 reply->reply.request_time = request->async->recv_time;
698
699 reply->listen = request->async->listen;
700 reply->packet_ctx = request->async->packet_ctx;
701
702 /*
703 * Update the various timers.
704 */
705 fr_time_elapsed_update(&worker->cpu_time, now, fr_time_add(now, reply->reply.processing_time));
706 fr_time_elapsed_update(&worker->wall_clock, reply->reply.request_time, now);
707
708 RDEBUG("Finished request");
709
710 /*
711 * Send the reply, which also polls the request queue.
712 */
713 if (fr_channel_send_reply(ch, reply) < 0) {
714 /*
715 * Should only happen if the TO_REQUESTOR
716 * channel is full, or it's not yet active.
717 *
718 * Not much we can do except complain
719 * loudly and cleanup the request.
720 */
721 RPERROR("Failed sending reply to network thread");
722 }
723
724 worker->stats.out++;
725
726 fr_assert(!fr_timer_armed(request->timeout));
727 fr_assert(!fr_heap_entry_inserted(request->runnable));
728
729 fr_dlist_entry_unlink(&request->listen_entry);
730
731#ifndef NDEBUG
732 request->async->el = NULL;
733 request->async->channel = NULL;
734 request->async->packet_ctx = NULL;
735 request->async->listen = NULL;
736#endif
737}
738
739/*
740 * talloc_typed_asprintf() is horrifically slow for printing
741 * simple numbers.
742 */
743static char *itoa_internal(TALLOC_CTX *ctx, uint64_t number)
744{
745 char buffer[32];
746 char *p;
747 char const *numbers = "0123456789";
748
749 p = buffer + 30;
750 *(p--) = '\0';
751
752 while (number > 0) {
753 *(p--) = numbers[number % 10];
754 number /= 10;
755 }
756
757 if (p[1]) return talloc_strdup(ctx, p + 1);
758
759 return talloc_strdup(ctx, "0");
760}
761
762/** Initialize various request fields needed by the worker.
763 *
764 */
765static inline CC_HINT(always_inline)
767{
768 /*
769 * For internal requests request->packet
770 * and request->reply are already populated.
771 */
772 if (!request->packet) MEM(request->packet = fr_packet_alloc(request, false));
773 if (!request->reply) MEM(request->reply = fr_packet_alloc(request, false));
774
775 request->packet->timestamp = now;
776 request->async = talloc_zero(request, fr_async_t);
777 request->async->request = request;
778 request->async->recv_time = now;
779 request->async->el = worker->el;
780 fr_dlist_entry_init(&request->async->entry);
781}
782
783static inline CC_HINT(always_inline)
785{
786 request->number = atomic_fetch_add_explicit(&request_number, 1, memory_order_seq_cst);
787 if (request->name) talloc_const_free(request->name);
788 request->name = itoa_internal(request, request->number);
789}
790
791static inline CC_HINT(always_inline)
793{
794 return fr_timer_list_num_events(worker->timeout);
795}
796
797static int _worker_request_deinit(request_t *request, UNUSED void *uctx)
798{
799 return request_slab_deinit(request);
800}
801
803{
804 fr_worker_t *worker = wc->worker;
805 int ret = -1;
806 request_t *request;
807 fr_listen_t *listen = cd->listen;
808
809 if (worker_num_requests(worker) >= (uint32_t) worker->config.max_requests) {
810 RATE_LIMIT_GLOBAL(ERROR, "Worker at max requests");
811 goto nak;
812 }
813
814 /*
815 * Receive a message to the worker queue, and decode it
816 * to a request.
817 */
818 fr_assert(listen != NULL);
819
820 request = request_slab_reserve(worker->slab);
821 if (!request) {
822 RATE_LIMIT_GLOBAL(ERROR, "Worker failed allocating new request");
823 goto nak;
824 }
825 /*
826 * Ensures that both the deinit function runs AND
827 * the request is returned to the slab if something
828 * calls talloc_free() on it.
829 */
830 request_slab_element_set_destructor(request, _worker_request_deinit, worker);
831
832 /*
833 * Have to initialise the request manually because namspace
834 * changes based on the listener that allocated it.
835 */
837 (&(request_init_args_t){ .namespace = listen->dict })) < 0) {
838 request_slab_release(request);
839 goto nak;
840 }
841
842 /*
843 * Do normal worker init that's shared between internal
844 * and external requests.
845 */
846 worker_request_init(worker, request, now);
848
849 /*
850 * Associate our interpreter with the request
851 */
852 unlang_interpret_set(request, worker->intp);
853
854 request->packet->timestamp = cd->request.recv_time; /* Legacy - Remove once everything looks at request->async */
855
856 /*
857 * Update the transport-specific fields.
858 */
859 request->async->channel = cd->channel.ch;
860
861 request->async->recv_time = cd->request.recv_time;
862
863 request->async->listen = listen;
864 request->async->packet_ctx = cd->packet_ctx;
865 request->priority = cd->priority;
866
867 /*
868 * Now that the "request" structure has been initialized, go decode the packet.
869 *
870 * Note that this also sets the "async process" function.
871 */
872 if (listen->app->decode) {
873 ret = listen->app->decode(listen->app_instance, request, cd->m.data, cd->m.data_size);
874 } else if (listen->app_io->decode) {
875 ret = listen->app_io->decode(listen->app_io_instance, request, cd->m.data, cd->m.data_size);
876 }
877
878 if (ret < 0) {
879 fail:
880 fr_assert(talloc_parent(request->stack) == request);
881 request_slab_release(request);
882
883 nak:
884 worker_nak(worker, cd, now);
885 return;
886 }
887
888 /*
889 * Set the entry point for this virtual server.
890 */
891 if (unlang_call_push(NULL, request, cd->listen->server_cs, UNLANG_TOP_FRAME) < 0) {
892 RERROR("Protocol failed to set 'process' function");
893 goto fail;
894 }
895
896 /*
897 * Look for conflicting / duplicate packets, but only if
898 * requested to do so.
899 */
900 if (request->async->listen->track_duplicates) {
901 request_t *old;
902
903 fr_rb_find((void **)&old, worker->dedup, request);
904 if (!old) {
905 goto insert_new;
906 }
907
908 fr_assert(old->async->listen == request->async->listen);
909 fr_assert(old->async->channel == request->async->channel);
910
911 /*
912 * There's a new packet. Do we keep the old one,
913 * or the new one? This decision is made by
914 * checking the recv_time, which is a
915 * nanosecond-resolution timer. If the time is
916 * identical, then the new packet is the same as
917 * the old one.
918 *
919 * If the new packet is a duplicate of the old
920 * one, then we can just discard the new one. We
921 * have to tell the channel that we've "eaten"
922 * this reply, so the sequence number should
923 * increase.
924 *
925 * @todo - fix the channel code to do queue
926 * depth, and not sequence / ack.
927 */
928 if (fr_time_eq(old->async->recv_time, request->async->recv_time)) {
929 RWARN("Discarding duplicate of request (%"PRIu64")", old->number);
930
931 fr_channel_null_reply(request->async->channel);
932 request_slab_release(request);
933
934 /*
935 * Signal there's a dup, and ignore the
936 * return code. We don't bother replying
937 * here, as an FD event or timer will
938 * wake up the request, and cause it to
939 * continue.
940 *
941 * @todo - the old request is NOT
942 * running, but is yielded. It MAY clean
943 * itself up, or do something...
944 */
946 worker->stats.dup++;
947
948 fr_message_done(&cd->m);
949 return;
950 }
951
952 /*
953 * Stop the old request, and decrement the number
954 * of active requests.
955 */
956 RWARN("Got conflicting packet for request (%" PRIu64 "), telling old request to stop", old->number);
957
959 worker->stats.dropped++;
960 (void) fr_rb_remove(NULL, worker->dedup, old); /* remove, but do NOT free it */
961
962 insert_new:
963 (void) fr_rb_insert(worker->dedup, request);
964 }
965
966 if (worker_request_time_tracking_start(worker, request, now) < 0) {
967 if (request->async->listen->track_duplicates) (void) fr_rb_remove(NULL, worker->dedup, request);
968 goto fail;
969 }
970
971 /*
972 * We're done with this message.
973 */
974 fr_message_done(&cd->m);
975
976 {
978
979 fr_rb_find((void **)&wl, worker->listeners, &(fr_worker_listen_t) { .listener = listen });
980 if (!wl) {
981 MEM(wl = talloc_zero(worker, fr_worker_listen_t));
982 fr_dlist_init(&wl->dlist, request_t, listen_entry);
983 wl->listener = listen;
984
985 (void) fr_rb_insert(worker->listeners, wl);
986 }
987
988 fr_dlist_insert_tail(&wl->dlist, request);
989 }
990
991 /*
992 * Track this request against the channel it came in on so
993 * worker_requests_cancel() has something to walk when the
994 * network signals CHANNEL_CLOSE.
995 */
996 fr_dlist_insert_tail(&wc->dlist, request->async);
997}
998
999/**
1000 * Track a request_t in the "runnable" heap.
1001 * Higher priorities take precedence, followed by lower sequence numbers
1002 */
1003static fr_cmp_ret_t worker_runnable_cmp(void const *one, void const *two)
1004{
1005 request_t const *a = one, *b = two;
1006 int ret;
1007
1008 /*
1009 * Prefer higher priority packets.
1010 */
1011 ret = CMP_PREFER_LARGER(a->priority, b->priority);
1012 if (ret != 0) return ret;
1013
1014 /*
1015 * Prefer packets which are further along in their processing sequence.
1016 */
1017 ret = CMP_PREFER_LARGER(a->sequence, b->sequence);
1018 if (ret != 0) return ret;
1019
1020 /*
1021 * Smaller timestamp (i.e. earlier) is more important.
1022 */
1023 return fr_time_cmp(a->async->recv_time, b->async->recv_time);
1024}
1025
1026/**
1027 * Track a request_t in the "dedup" tree
1028 */
1029static fr_cmp_ret_t worker_dedup_cmp(void const *one, void const *two)
1030{
1031 int ret;
1032 request_t const *a = one, *b = two;
1033
1034 ret = CMP(a->async->listen, b->async->listen);
1035 if (ret) return ret;
1036
1037 return CMP(a->async->packet_ctx, b->async->packet_ctx);
1038}
1039
1040/** Destroy a worker
1041 *
1042 * The input channels are signaled, and local messages are cleaned up.
1043 *
1044 * This should be called to _EXPLICITLY_ destroy a worker, when some fatal
1045 * error has occurred on the worker side, and we need to destroy it.
1046 *
1047 * We signal all pending requests in the backlog to stop, and tell the
1048 * network side that it should not send us any more requests.
1049 *
1050 * @param[in] worker the worker to destroy.
1051 */
1053{
1054 int i, count, ret;
1055
1056// WORKER_VERIFY;
1057
1058 /*
1059 * Stop any new requests running with this interpreter
1060 */
1062
1063 /*
1064 * Destroy all of the active requests. These are ones
1065 * which are still waiting for timers or file descriptor
1066 * events.
1067 */
1068 count = 0;
1069
1070 /*
1071 * Force the timeout event to fire for all requests that
1072 * are still running.
1073 */
1074 ret = fr_timer_list_force_run(worker->timeout);
1075 if (unlikely(ret < 0)) {
1076 fr_assert_msg(0, "Failed to force run the timeout list");
1077 } else {
1078 count += ret;
1079 }
1080
1082
1083 DEBUG("Worker is exiting - stopped %u requests", count);
1084
1085 /*
1086 * Signal the channels that we're closing.
1087 *
1088 * The other end owns the channel, and will take care of
1089 * popping messages in the TO_RESPONDER queue, and marking
1090 * them FR_MESSAGE_DONE. It will ignore the messages in
1091 * the TO_REQUESTOR queue, as we own those. They will be
1092 * automatically freed when our talloc context is freed.
1093 */
1094 for (i = 0; i < worker->config.max_channels; i++) {
1095 if (!worker->channel[i].ch) continue;
1096
1097 worker_requests_cancel(&worker->channel[i]);
1098
1099 fr_assert_msg(fr_dlist_num_elements(&worker->channel[i].dlist) == 0,
1100 "Pending messages in channel after cancelling request");
1101
1103 }
1104
1105 talloc_free(worker);
1106}
1107
1108/** Internal request (i.e. one generated by the interpreter) is now complete
1109 *
1110 */
1111static void _worker_request_internal_init(request_t *request, void *uctx)
1112{
1113 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1114 fr_time_t now = fr_time();
1115
1116 worker_request_init(worker, request, now);
1117
1118 /*
1119 * Requests generated by the interpreter
1120 * are always marked up as internal.
1121 */
1123 if (worker_request_time_tracking_start(worker, request, now) < 0) {
1125 }
1126}
1127
1128
1129/** External request is now complete
1130 *
1131 */
1132static void _worker_request_done_external(request_t *request, UNUSED rlm_rcode_t rcode, void *uctx)
1133{
1134 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1135 fr_time_t now = fr_time();
1136
1137 /*
1138 * All external requests MUST have a listener.
1139 */
1141 fr_assert(request->async->listen != NULL);
1142
1143 /*
1144 * Only real packets are in the dedup tree. And even
1145 * then, only some of the time.
1146 */
1147 if (request->async->listen->track_duplicates && fr_rb_node_inline_in_tree(&request->dedup_node)) {
1148 (void) fr_rb_delete(worker->dedup, request);
1149 }
1150
1151 /*
1152 * If we're running a real request, then the final
1153 * indentation MUST be zero. Otherwise we skipped
1154 * something!
1155 *
1156 * Also check that the request is NOT marked as
1157 * "yielded", but is in fact done.
1158 *
1159 * @todo - check that the stack is at frame 0, otherwise
1160 * more things have gone wrong.
1161 */
1162 fr_assert_msg(request_is_internal(request) || request_is_detached(request) || (request->log.indent.unlang == 0),
1163 "Request %s bad log indentation - expected 0 got %u", request->name, request->log.indent.unlang);
1165 "Request %s is marked as yielded at end of processing", request->name);
1167 "Request %s stack depth %u > 0", request->name, unlang_interpret_stack_depth(request));
1168 RDEBUG("Done request");
1169
1170 /*
1171 * The request is done. Track that.
1172 */
1173 worker_request_time_tracking_end(worker, request, now);
1174
1175 /*
1176 * Remove it from the list of requests associated with this channel.
1177 */
1178 if (fr_dlist_entry_in_list(&request->async->entry)) {
1179 fr_worker_channel_t *wc = fr_channel_responder_uctx_get(request->async->channel);
1180
1181 fr_dlist_remove(&wc->dlist, request->async);
1182 }
1183
1184 /*
1185 * These conditions are true when the server is
1186 * exiting and we're stopping all the requests.
1187 *
1188 * This should never happen otherwise.
1189 */
1190 if (unlikely(!fr_channel_active(request->async->channel))) {
1191 fr_dlist_entry_unlink(&request->listen_entry);
1192 request_slab_release(request);
1193 return;
1194 }
1195
1196 worker_send_reply(worker, request, !unlang_request_is_cancelled(request), now);
1197 request_slab_release(request);
1198}
1199
1200/** Internal request (i.e. one generated by the interpreter) is now complete
1201 *
1202 * Whatever generated the request is now responsible for freeing it.
1203 */
1204static void _worker_request_done_internal(request_t *request, UNUSED rlm_rcode_t rcode, void *uctx)
1205{
1206 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1207
1208 worker_request_time_tracking_end(worker, request, fr_time());
1209
1210 fr_assert(!fr_heap_entry_inserted(request->runnable));
1211 fr_assert(!fr_timer_armed(request->timeout));
1212 fr_assert(!fr_dlist_entry_in_list(&request->async->entry));
1213}
1214
1215/** Detached request (i.e. one generated by the interpreter with no parent) is now complete
1216 *
1217 * As the request has no parent, then there's nothing to free it
1218 * so we have to.
1219 */
1220static void _worker_request_done_detached(request_t *request, UNUSED rlm_rcode_t rcode, UNUSED void *uctx)
1221{
1222 /*
1223 * No time tracking for detached requests
1224 * so we don't need to call
1225 * worker_request_time_tracking_end.
1226 */
1227 fr_assert(!fr_heap_entry_inserted(request->runnable));
1228
1229 /*
1230 * Normally worker_request_time_tracking_end
1231 * would remove the request from the time
1232 * order heap, but we need to do that for
1233 * detached requests.
1234 */
1235 TALLOC_FREE(request->timeout);
1236
1237 fr_assert(!fr_dlist_entry_in_list(&request->async->entry));
1238
1239 /*
1240 * Detached requests have to be freed by us
1241 * as nothing else can free them.
1242 *
1243 * All other requests must be freed by the
1244 * code which allocated them.
1245 */
1246 talloc_free(request);
1247}
1248
1249
1250/** Make us responsible for running the request
1251 *
1252 */
1253static void _worker_request_detach(request_t *request, void *uctx)
1254{
1255 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1256 fr_time_t now = fr_time();
1257
1258 RDEBUG4("%s - Request detaching", __FUNCTION__);
1259
1260 if (request_is_detachable(request)) {
1261 /*
1262 * End the time tracking... We don't track detached requests,
1263 * because they don't contribute for the time consumed by an
1264 * external request.
1265 */
1266 if (request->async->tracking.state == FR_TIME_TRACKING_YIELDED) {
1267 RDEBUG3("Forcing time tracking to running state, from yielded, for request detach");
1268 fr_time_tracking_resume(&request->async->tracking, now);
1269 }
1270 worker_request_time_tracking_end(worker, request, now);
1271
1272 if (request_detach(request) < 0) RPEDEBUG("Failed detaching request");
1273
1274 RDEBUG3("Request is detached");
1275 } else {
1276 fr_assert_msg(0, "Request is not detachable");
1277 }
1278
1279 return;
1280}
1281
1282/** Request is now runnable
1283 *
1284 */
1285static void _worker_request_runnable(request_t *request, void *uctx)
1286{
1287 fr_worker_t *worker = uctx;
1288
1289 RDEBUG4("%s - Request marked as runnable", __FUNCTION__);
1290 fr_heap_insert(&worker->runnable, request);
1291}
1292
1293/** Interpreter yielded request
1294 *
1295 */
1296static void _worker_request_yield(request_t *request, UNUSED void *uctx)
1297{
1298 RDEBUG4("%s - Request yielded", __FUNCTION__);
1299 if (likely(!request_is_detached(request))) fr_time_tracking_yield(&request->async->tracking, fr_time());
1300}
1301
1302/** Interpreter is starting to work on request again
1303 *
1304 */
1305static void _worker_request_resume(request_t *request, UNUSED void *uctx)
1306{
1307 RDEBUG4("%s - Request resuming", __FUNCTION__);
1308 if (likely(!request_is_detached(request))) fr_time_tracking_resume(&request->async->tracking, fr_time());
1309}
1310
1311/** Check if a request is scheduled
1312 *
1313 */
1314static bool _worker_request_scheduled(request_t const *request, UNUSED void *uctx)
1315{
1316 return fr_heap_entry_inserted(request->runnable);
1317}
1318
1319/** Update a request's priority
1320 *
1321 */
1322static void _worker_request_prioritise(request_t *request, void *uctx)
1323{
1324 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1325
1326 RDEBUG4("%s - Request priority changed", __FUNCTION__);
1327
1328 /* Extract the request from the runnable queue _if_ it's in the runnable queue */
1329 if (fr_heap_extract(&worker->runnable, request) < 0) return;
1330
1331 /* Reinsert it to re-evaluate its new priority */
1332 fr_heap_insert(&worker->runnable, request);
1333}
1334
1335/** Run a request
1336 *
1337 * Until it either yields, or is done.
1338 *
1339 * This function is also responsible for sending replies, and
1340 * cleaning up the request.
1341 *
1342 * @param[in] worker the worker
1343 * @param[in] start the current time
1344 */
1345static inline CC_HINT(always_inline) void worker_run_request(fr_worker_t *worker, fr_time_t start)
1346{
1347 request_t *request;
1348 fr_time_t now;
1349
1351
1352 now = start;
1353
1354 /*
1355 * Busy-loop running requests for 1ms. We still poll the
1356 * event loop 1000 times a second, OR when there's no
1357 * more work to do. This allows us to make progress with
1358 * ongoing requests, at the expense of sometimes ignoring
1359 * new ones.
1360 */
1361 while (fr_time_delta_lt(fr_time_sub(now, start), fr_time_delta_from_msec(1)) &&
1362 (fr_heap_pop((void **)&request, &worker->runnable) == 0) && request) {
1363
1364 REQUEST_VERIFY(request);
1365 fr_assert(!fr_heap_entry_inserted(request->runnable));
1366
1367 /*
1368 * For real requests, if the channel is gone,
1369 * just stop the request and free it.
1370 */
1371 if (request->async->channel && !fr_channel_active(request->async->channel)) {
1372 worker_stop_request(request);
1373 continue;
1374 }
1375
1377
1378 now = fr_time();
1379 }
1380}
1381
1382/** Create a worker
1383 *
1384 * @param[in] ctx the talloc context
1385 * @param[in] name the name of this worker
1386 * @param[in] el the event list
1387 * @param[in] logger the destination for all logging messages
1388 * @param[in] lvl log level
1389 * @param[in] config various configuration parameters
1390 * @return
1391 * - NULL on error
1392 * - fr_worker_t on success
1393 */
1394fr_worker_t *fr_worker_alloc(TALLOC_CTX *ctx, fr_event_list_t *el, char const *name, fr_log_t const *logger, fr_log_lvl_t lvl,
1396{
1397 fr_worker_t *worker;
1398
1399 worker = talloc_zero(ctx, fr_worker_t);
1400 if (!worker) {
1401nomem:
1402 fr_strerror_const("Failed allocating memory");
1403 return NULL;
1404 }
1405
1406 worker->name = talloc_strdup(worker, name); /* thread locality */
1407
1408 if (config) worker->config = *config;
1409
1410#define CHECK_CONFIG(_x, _min, _max) do { \
1411 if (!worker->config._x) worker->config._x = _min; \
1412 if (worker->config._x < _min) worker->config._x = _min; \
1413 if (worker->config._x > _max) worker->config._x = _max; \
1414 } while (0)
1415
1416#define CHECK_CONFIG_TIME_DELTA(_x, _min, _max) do { \
1417 if (fr_time_delta_lt(worker->config._x, _min)) worker->config._x = _min; \
1418 if (fr_time_delta_gt(worker->config._x, _max)) worker->config._x = _max; \
1419 } while (0)
1420
1421 CHECK_CONFIG(max_requests,1024,(1 << 30));
1422 CHECK_CONFIG(max_channels, 64, 1024);
1423 CHECK_CONFIG(reuse.child_pool_size, 4096, 65536);
1424 CHECK_CONFIG(message_set_size, 1024, 8192);
1425 CHECK_CONFIG(ring_buffer_size, (1 << 17), (1 << 20));
1427
1428 worker->channel = talloc_zero_array(worker, fr_worker_channel_t, worker->config.max_channels);
1429 if (!worker->channel) {
1430 talloc_free(worker);
1431 goto nomem;
1432 }
1433
1434 worker->thread_id = pthread_self();
1435 worker->el = el;
1436 worker->log = logger;
1437 worker->lvl = lvl;
1438
1439 /*
1440 * The worker thread starts now. Manually initialize it,
1441 * because we're tracking request time, not the time that
1442 * the worker thread is running.
1443 */
1444 memset(&worker->tracking, 0, sizeof(worker->tracking));
1445
1446 worker->aq_control = fr_atomic_queue_talloc(worker, 1024);
1447 if (!worker->aq_control) {
1448 fr_strerror_const("Failed creating atomic queue");
1449 fail:
1450 talloc_free(worker);
1451 return NULL;
1452 }
1453
1454 worker->control = fr_control_create(worker, el, worker->aq_control, 7);
1455 if (!worker->control) {
1456 fr_strerror_const_push("Failed creating control plane");
1457 goto fail;
1458 }
1459
1461 fr_strerror_const_push("Failed adding control channel");
1462 goto fail;
1463 }
1464
1466 fr_strerror_const_push("Failed adding callback for listeners");
1467 goto fail;
1468 }
1469
1470 if (fr_control_open(worker->control) < 0) {
1471 fr_strerror_const_push("Failed opening control plane");
1472 goto fail;
1473 }
1474
1475 worker->runnable = fr_heap_talloc_alloc(worker, worker_runnable_cmp, request_t, runnable, 0);
1476 if (!worker->runnable) {
1477 fr_strerror_const("Failed creating runnable heap");
1478 goto fail;
1479 }
1480
1481 worker->timeout = fr_timer_list_ordered_alloc(worker, el->tl);
1482 if (!worker->timeout) {
1483 fr_strerror_const("Failed creating timeouts list");
1484 goto fail;
1485 }
1486
1487 worker->dedup = fr_rb_inline_talloc_alloc(worker, request_t, dedup_node, worker_dedup_cmp, NULL);
1488 if (!worker->dedup) {
1489 fr_strerror_const("Failed creating de_dup tree");
1490 goto fail;
1491 }
1492
1494 if (!worker->listeners) {
1495 fr_strerror_const("Failed creating listener tree");
1496 goto fail;
1497 }
1498
1499 worker->intp = unlang_interpret_init(worker, el,
1501 .init_internal = _worker_request_internal_init,
1502
1503 .done_external = _worker_request_done_external,
1504 .done_internal = _worker_request_done_internal,
1505 .done_detached = _worker_request_done_detached,
1506
1507 .detach = _worker_request_detach,
1508 .yield = _worker_request_yield,
1509 .resume = _worker_request_resume,
1510 .mark_runnable = _worker_request_runnable,
1511
1512 .scheduled = _worker_request_scheduled,
1513 .prioritise = _worker_request_prioritise
1514 },
1515 worker);
1516 if (!worker->intp){
1517 fr_strerror_const("Failed initialising interpreter");
1518 goto fail;
1519 }
1520
1521 {
1524
1525 if (!(worker->slab = request_slab_list_alloc(worker, el, &worker->config.reuse, NULL, NULL,
1526 UNCONST(void *, worker), true, false))) {
1527 fr_strerror_const("Failed creating request slab list");
1528 goto fail;
1529 }
1530 }
1531
1533
1534 return worker;
1535}
1536
1537
1538/** The main loop and entry point of the stand-alone worker thread.
1539 *
1540 * Where there is only one thread, the event loop runs fr_worker_pre_event() and fr_worker_post_event()
1541 * instead, And then fr_worker_post_event() takes care of calling worker_run_request() to actually run the
1542 * request.
1543 *
1544 * @param[in] worker the worker data structure to manage
1545 */
1547{
1549
1550 while (true) {
1551 bool wait_for_event;
1552 int num_events;
1553
1555
1556 /*
1557 * There are runnable requests. We still service
1558 * the event loop, but we don't wait for events.
1559 */
1560 wait_for_event = (fr_heap_num_elements(worker->runnable) == 0);
1561 if (wait_for_event) {
1562 if (worker->exiting && (worker_num_requests(worker) == 0)) break;
1563
1564 DEBUG4("Ready to process requests");
1565 }
1566
1567 /*
1568 * Check the event list. If there's an error
1569 * (e.g. exit), we stop looping and clean up.
1570 */
1571 DEBUG4("Gathering events - %s", wait_for_event ? "will wait" : "Will not wait");
1572 num_events = fr_event_corral(worker->el, fr_time(), wait_for_event);
1573 if (num_events < 0) {
1574 if (fr_event_loop_exiting(worker->el)) {
1575 DEBUG4("Event loop exiting");
1576 break;
1577 }
1578
1579 PERROR("Failed retrieving events");
1580 break;
1581 }
1582
1583 DEBUG4("%u event(s) pending", num_events);
1584
1585 /*
1586 * Service outstanding events.
1587 */
1588 if (num_events > 0) {
1589 DEBUG4("Servicing event(s)");
1590 fr_event_service(worker->el);
1591 }
1592
1593 /*
1594 * Run any outstanding requests.
1595 */
1596 worker_run_request(worker, fr_time());
1597 }
1598}
1599
1600/** Pre-event handler
1601 *
1602 * This should be run ONLY in single-threaded mode!
1603 */
1605{
1606 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1607 request_t *request;
1608
1609 request = fr_heap_peek(worker->runnable);
1610 if (!request) return 0;
1611
1612 /*
1613 * There's work to do. Tell the event handler to poll
1614 * for IO / timers, but also immediately return to the
1615 * calling function, which has more work to do.
1616 */
1617 return 1;
1618}
1619
1620
1621/** Post-event handler
1622 *
1623 * This should be run ONLY in single-threaded mode!
1624 */
1626{
1627 fr_worker_t *worker = talloc_get_type_abort(uctx, fr_worker_t);
1628
1629 worker_run_request(worker, fr_time()); /* Event loop time can be too old, and trigger asserts */
1630}
1631
1632/** Print debug information about the worker structure
1633 *
1634 * @param[in] worker the worker
1635 * @param[in] fp the file where the debug output is printed.
1636 */
1637void fr_worker_debug(fr_worker_t *worker, FILE *fp)
1638{
1640
1641 fprintf(fp, "\tnum_channels = %d\n", worker->num_channels);
1642 fprintf(fp, "\tstats.in = %" PRIu64 "\n", worker->stats.in);
1643
1644 fprintf(fp, "\tcalculated (predicted) total CPU time = %" PRIu64 "\n",
1645 fr_time_delta_unwrap(worker->predicted) * worker->stats.in);
1646 if (worker->stats.in) {
1647 fprintf(fp, "\tcalculated (counted) per request time = %" PRIu64 "\n",
1649 }
1650
1651 fr_time_tracking_debug(&worker->tracking, fp);
1652
1653}
1654
1655/** Create a channel to the worker
1656 *
1657 * Called by the master (i.e. network) thread when it needs to create
1658 * a new channel to a particuler worker.
1659 *
1660 * @param[in] worker the worker
1661 * @param[in] master the control plane of the master
1662 * @param[in] ctx the context in which the channel will be created
1663 */
1665{
1666 fr_channel_t *ch;
1667 pthread_t id;
1668 bool same;
1669
1671
1672 id = pthread_self();
1673 same = (pthread_equal(id, worker->thread_id) != 0);
1674
1675 ch = fr_channel_create(ctx, master, worker->control, same);
1676 if (!ch) return NULL;
1677
1678
1679 /*
1680 * Tell the worker about the channel
1681 */
1682 if (fr_channel_signal_open(ch) < 0) {
1683 talloc_free(ch);
1684 return NULL;
1685 }
1686
1687 return ch;
1688}
1689
1691{
1692 fr_ring_buffer_t *rb;
1693
1694 /*
1695 * Skip a bunch of work if we're already in the worker thread.
1696 */
1697 if (is_worker_thread(worker)) {
1698 return fr_worker_listen_cancel_self(worker, li);
1699 }
1700
1701 rb = fr_worker_rb_init();
1702 if (!rb) return -1;
1703
1704 return fr_control_message_send(worker->control, rb, FR_CONTROL_ID_LISTEN_DEAD, &li, sizeof(li));
1705}
1706
1707#ifdef WITH_VERIFY_PTR
1708/** Verify the worker data structures.
1709 *
1710 * @param[in] worker the worker
1711 */
1712static void worker_verify(fr_worker_t *worker)
1713{
1714 int i;
1715
1716 (void) talloc_get_type_abort(worker, fr_worker_t);
1717 fr_atomic_queue_verify(worker->aq_control);
1718
1719 fr_assert(worker->control != NULL);
1720 (void) talloc_get_type_abort(worker->control, fr_control_t);
1721
1722 fr_assert(worker->el != NULL);
1723 (void) talloc_get_type_abort(worker->el, fr_event_list_t);
1724
1725 fr_assert(worker->runnable != NULL);
1726 (void) talloc_get_type_abort(worker->runnable, fr_heap_t);
1727
1728 fr_assert(worker->dedup != NULL);
1729 (void) talloc_get_type_abort(worker->dedup, fr_rb_tree_t);
1730
1731 for (i = 0; i < worker->config.max_channels; i++) {
1732 if (!worker->channel[i].ch) continue;
1733
1734 (void) talloc_get_type_abort(worker->channel[i].ch, fr_channel_t);
1735 }
1736}
1737#endif
1738
1739int fr_worker_stats(fr_worker_t const *worker, int num, uint64_t *stats)
1740{
1741 if (num < 0) return -1;
1742 if (num == 0) return 0;
1743
1744 stats[0] = worker->stats.in;
1745 if (num >= 2) stats[1] = worker->stats.out;
1746 if (num >= 3) stats[2] = worker->stats.dup;
1747 if (num >= 4) stats[3] = worker->stats.dropped;
1748 if (num >= 5) stats[4] = worker->num_naks;
1749 if (num >= 6) stats[5] = worker->num_active;
1750
1751 if (num <= 6) return num;
1752
1753 return 6;
1754}
1755
1756static int cmd_stats_worker(FILE *fp, UNUSED FILE *fp_err, void *ctx, fr_cmd_info_t const *info)
1757{
1758 fr_worker_t const *worker = ctx;
1759 fr_time_delta_t when;
1760
1761 if ((info->argc == 0) || (strcmp(info->argv[0], "count") == 0)) {
1762 fprintf(fp, "count.in\t\t\t%" PRIu64 "\n", worker->stats.in);
1763 fprintf(fp, "count.out\t\t\t%" PRIu64 "\n", worker->stats.out);
1764 fprintf(fp, "count.dup\t\t\t%" PRIu64 "\n", worker->stats.dup);
1765 fprintf(fp, "count.dropped\t\t\t%" PRIu64 "\n", worker->stats.dropped);
1766 fprintf(fp, "count.naks\t\t\t%" PRIu64 "\n", worker->num_naks);
1767 fprintf(fp, "count.active\t\t\t%" PRIu64 "\n", worker->num_active);
1768 fprintf(fp, "count.runnable\t\t\t%u\n", fr_heap_num_elements(worker->runnable));
1769 }
1770
1771 if ((info->argc == 0) || (strcmp(info->argv[0], "cpu") == 0)) {
1772 when = worker->predicted;
1773 fprintf(fp, "cpu.request_time_rtt\t\t%.9f\n", fr_time_delta_unwrap(when) / (double)NSEC);
1774
1775 when = worker->tracking.running_total;
1776 if (fr_time_delta_ispos(when) && (worker->stats.in > worker->stats.dropped)) {
1777 when = fr_time_delta_div(when, fr_time_delta_wrap(worker->stats.in - worker->stats.dropped));
1778 }
1779 fprintf(fp, "cpu.average_request_time\t%.9f\n", fr_time_delta_unwrap(when) / (double)NSEC);
1780
1781 when = worker->tracking.running_total;
1782 fprintf(fp, "cpu.used\t\t\t%.6f\n", fr_time_delta_unwrap(when) / (double)NSEC);
1783
1784 when = worker->tracking.waiting_total;
1785 fprintf(fp, "cpu.waiting\t\t\t%.3f\n", fr_time_delta_unwrap(when) / (double)NSEC);
1786
1787 fr_time_elapsed_fprint(fp, &worker->cpu_time, "cpu.requests", 4);
1788 fr_time_elapsed_fprint(fp, &worker->wall_clock, "time.requests", 4);
1789 }
1790
1791 return 0;
1792}
1793
1795 {
1796 .parent = "stats",
1797 .name = "worker",
1798 .help = "Statistics for workers threads.",
1799 .read_only = true
1800 },
1801
1802 {
1803 .parent = "stats worker",
1804 .add_name = true,
1805 .name = "self",
1806 .syntax = "[(count|cpu)]",
1807 .func = cmd_stats_worker,
1808 .help = "Show statistics for a specific worker thread.",
1809 .read_only = true
1810 },
1811
1813};
static int const char char buffer[256]
Definition acutest.h:576
fr_io_encode_t encode
Pack fr_pair_ts back into a byte array.
Definition app_io.h:55
size_t default_reply_size
same for replies
Definition app_io.h:40
size_t default_message_size
Usually maximum message size.
Definition app_io.h:39
fr_io_nak_t nak
Function to send a NAK.
Definition app_io.h:62
fr_io_decode_t decode
Translate raw bytes into fr_pair_ts and metadata.
Definition app_io.h:54
fr_io_decode_t decode
Translate raw bytes into fr_pair_ts and metadata.
Definition application.h:80
fr_io_encode_t encode
Pack fr_pair_ts back into a byte array.
Definition application.h:85
#define _Thread_local
Definition atexit.h:213
#define fr_atexit_thread_local(_name, _free, _uctx)
Definition atexit.h:224
fr_atomic_queue_t * fr_atomic_queue_talloc(TALLOC_CTX *ctx, size_t size)
Create fixed-size atomic queue.
Structure to hold the atomic queue.
#define UNCONST(_type, _ptr)
Remove const qualification from a pointer.
Definition build.h:186
#define RCSID(id)
Definition build.h:560
#define NDEBUG_UNUSED
Definition build.h:395
#define CMP_PREFER_LARGER(_a, _b)
Evaluates to -1 for a > b, and +1 for a < b.
Definition build.h:109
#define CMP(_a, _b)
Same as CMP_PREFER_SMALLER use when you don't really care about ordering, you just want an ordering.
Definition build.h:113
#define unlikely(_x)
Definition build.h:455
#define UNUSED
Definition build.h:384
unlang_action_t unlang_call_push(unlang_result_t *p_result, request_t *request, CONF_SECTION *server_cs, bool top_frame)
Push a virtual server CONF_SECTION as a call frame onto the stack.
Definition call.c:151
fr_table_num_sorted_t const channel_signals[]
Definition channel.c:151
unsigned int fr_channel_responder_discard(fr_channel_t *ch)
Discard any requests the requestor queued but we never received.
Definition channel.c:874
fr_channel_t * fr_channel_create(TALLOC_CTX *ctx, fr_control_t *requestor, fr_control_t *responder, bool same)
Create a new channel.
Definition channel.c:181
fr_channel_event_t fr_channel_service_message(fr_time_t when, fr_channel_t **p_channel, void const *data, size_t data_size)
Service a control-plane message.
Definition channel.c:687
void * fr_channel_responder_uctx_get(fr_channel_t *ch)
Get responder-specific data from a channel.
Definition channel.c:936
bool fr_channel_recv_request(fr_channel_t *ch)
Receive a request message from the channel.
Definition channel.c:470
int fr_channel_null_reply(fr_channel_t *ch)
Don't send a reply message into the channel.
Definition channel.c:626
void fr_channel_responder_uctx_add(fr_channel_t *ch, void *uctx)
Add responder-specific data to a channel.
Definition channel.c:924
int fr_channel_set_recv_request(fr_channel_t *ch, fr_channel_recv_callback_t recv_request, void *uctx)
Definition channel.c:977
int fr_channel_send_reply(fr_channel_t *ch, fr_channel_data_t *cd)
Send a reply message into the channel.
Definition channel.c:509
bool fr_channel_active(fr_channel_t *ch)
Check if a channel is active.
Definition channel.c:829
int fr_channel_responder_ack_close(fr_channel_t *ch)
Acknowledge that the channel is closing.
Definition channel.c:895
int fr_channel_signal_open(fr_channel_t *ch)
Send a channel to a responder.
Definition channel.c:991
A full channel, which consists of two ends.
Definition channel.c:142
fr_message_t m
the message header
Definition channel.h:107
fr_channel_event_t
Definition channel.h:69
@ FR_CHANNEL_NOOP
Definition channel.h:76
@ FR_CHANNEL_EMPTY
Definition channel.h:77
@ FR_CHANNEL_CLOSE
Definition channel.h:74
@ FR_CHANNEL_ERROR
Definition channel.h:70
@ FR_CHANNEL_DATA_READY_REQUESTOR
Definition channel.h:72
@ FR_CHANNEL_OPEN
Definition channel.h:73
@ FR_CHANNEL_DATA_READY_RESPONDER
Definition channel.h:71
void * packet_ctx
Packet specific context for holding client information, and other proto_* specific information that n...
Definition channel.h:144
fr_listen_t * listen
for tracking packet transport, etc.
Definition channel.h:148
#define FR_CONTROL_ID_CHANNEL
Definition channel.h:67
uint32_t priority
Priority of this packet.
Definition channel.h:142
Channel information which is added to a message.
Definition channel.h:106
char const * parent
e.g. "show module"
Definition command.h:52
#define CMD_TABLE_END
Definition command.h:62
#define FR_CONTROL_MAX_SIZE
Definition control.h:51
#define FR_CONTROL_MAX_MESSAGES
Definition control.h:50
#define fr_cond_assert(_x)
Calls panic_action ifndef NDEBUG, else logs error and evaluates to value of _x.
Definition debug.h:172
#define fr_assert_msg(_x, _msg,...)
Calls panic_action ifndef NDEBUG, else logs error and causes the server to exit immediately with code...
Definition debug.h:243
#define fr_cond_assert_msg(_x, _fmt,...)
Calls panic_action ifndef NDEBUG, else logs error and evaluates to value of _x.
Definition debug.h:189
#define MEM(x)
Definition debug.h:38
#define ERROR(fmt,...)
Definition dhcpclient.c:40
#define DEBUG(fmt,...)
Definition dhcpclient.c:38
#define fr_dlist_init(_head, _type, _field)
Initialise the head structure of a doubly linked list.
Definition dlist.h:242
static void * fr_dlist_remove(fr_dlist_head_t *list_head, void *ptr)
Remove an item from the list.
Definition dlist.h:620
static bool fr_dlist_entry_in_list(fr_dlist_t const *entry)
Check if a list entry is part of a list.
Definition dlist.h:145
static void fr_dlist_entry_unlink(fr_dlist_t *entry)
Remove an item from the dlist when we don't have access to the head.
Definition dlist.h:128
static unsigned int fr_dlist_num_elements(fr_dlist_head_t const *head)
Return the number of elements in the dlist.
Definition dlist.h:921
static void * fr_dlist_pop_head(fr_dlist_head_t *list_head)
Remove the head item in a list.
Definition dlist.h:654
static int fr_dlist_insert_tail(fr_dlist_head_t *list_head, void *ptr)
Insert an item into the tail of a list.
Definition dlist.h:360
static void fr_dlist_entry_init(fr_dlist_t *entry)
Initialise a linked list without metadata.
Definition dlist.h:120
Head of a doubly linked list.
Definition dlist.h:51
int fr_heap_insert(fr_heap_t **hp, void *data)
Insert a new element into the heap.
Definition heap.c:149
int fr_heap_pop(void **out, fr_heap_t **hp)
Remove a node from the heap.
Definition heap.c:359
int fr_heap_extract(fr_heap_t **hp, void *data)
Remove a node from the heap.
Definition heap.c:259
static void * fr_heap_peek(fr_heap_t *h)
Return the item from the top of the heap but don't pop it.
Definition heap.h:138
static bool fr_heap_entry_inserted(fr_heap_index_t heap_idx)
Check if an entry is inserted into a heap.
Definition heap.h:126
static unsigned int fr_heap_num_elements(fr_heap_t *h)
Return the number of elements in the heap.
Definition heap.h:181
#define fr_heap_talloc_alloc(_ctx, _cmp, _talloc_type, _field, _init)
Creates a heap that verifies elements are of a specific talloc type.
Definition heap.h:117
The main heap structure.
Definition heap.h:68
talloc_free(hp)
rlm_rcode_t unlang_interpret(request_t *request, bool running)
Run the interpreter for a current request.
Definition interpret.c:1302
void unlang_interpret_set(request_t *request, unlang_interpret_t *intp)
Set a specific interpreter for a request.
Definition interpret.c:2519
int unlang_interpret_stack_depth(request_t *request)
Return the depth of the request's stack.
Definition interpret.c:1921
void unlang_interpret_set_thread_default(unlang_interpret_t *intp)
Set the default interpreter for this thread.
Definition interpret.c:2550
unlang_interpret_t * unlang_interpret_init(TALLOC_CTX *ctx, fr_event_list_t *el, unlang_request_func_t *funcs, void *uctx)
Initialize a unlang compiler / interpret.
Definition interpret.c:2478
bool unlang_request_is_cancelled(request_t const *request)
Return whether a request has been cancelled.
Definition interpret.c:1971
void unlang_interpret_signal(request_t *request, fr_signal_t action)
Send a signal (usually stop) to a request.
Definition interpret.c:1789
bool unlang_interpret_is_resumable(request_t *request)
Check if a request as resumable.
Definition interpret.c:1990
#define UNLANG_REQUEST_RESUME
Definition interpret.h:48
#define UNLANG_TOP_FRAME
Definition interpret.h:36
External functions provided by the owner of the interpret.
Definition interpret.h:116
uint64_t out
Definition base.h:43
uint64_t dup
Definition base.h:44
uint64_t dropped
Definition base.h:45
uint64_t in
Definition base.h:42
fr_control_t * fr_control_create(TALLOC_CTX *ctx, fr_event_list_t *el, fr_atomic_queue_t *aq, size_t num_callbacks)
Create a control-plane signaling path.
Definition control.c:152
int fr_control_open(fr_control_t *c)
Open the control-plane signalling path.
Definition control.c:176
int fr_control_message_send(fr_control_t *c, fr_ring_buffer_t *rb, uint32_t id, void *data, size_t data_size)
Send a control-plane message.
Definition control.c:355
int fr_control_callback_add(fr_control_t **c, uint32_t id, fr_control_callback_t callback, void *uctx)
Register a callback for an ID.
Definition control.c:444
The control structure.
Definition control.c:76
#define PERROR(_fmt,...)
Definition log.h:233
#define DEBUG3(_fmt,...)
Definition log.h:271
#define RDEBUG3(fmt,...)
Definition log.h:360
#define RWARN(fmt,...)
Definition log.h:314
#define PWARN(_fmt,...)
Definition log.h:232
#define RERROR(fmt,...)
Definition log.h:315
#define DEBUG4(_fmt,...)
Definition log.h:272
#define RPERROR(fmt,...)
Definition log.h:319
#define RPEDEBUG(fmt,...)
Definition log.h:393
#define RDEBUG4(fmt,...)
Definition log.h:361
#define RATE_LIMIT_GLOBAL(_log, _fmt,...)
Rate limit messages using a global limiting entry.
Definition log.h:658
void fr_event_service(fr_event_list_t *el)
Service any outstanding timer or file descriptor events.
Definition event.c:2204
int fr_event_corral(fr_event_list_t *el, fr_time_t now, bool wait)
Gather outstanding timer and file descriptor events.
Definition event.c:2072
int fr_event_post_delete(fr_event_list_t *el, fr_event_post_cb_t callback, void *uctx)
Delete a post-event callback from the event list.
Definition event.c:2050
#define fr_time()
Definition event.c:60
bool fr_event_loop_exiting(fr_event_list_t *el)
Check to see whether the event loop is in the process of exiting.
Definition event.c:2393
Stores all information relating to an event list.
Definition event.c:377
fr_log_lvl_t
Definition log.h:64
fr_packet_t * fr_packet_alloc(TALLOC_CTX *ctx, bool new_vector)
Allocate a new fr_packet_t.
Definition packet.c:38
request_t * request
back-pointer to the owning request so anything that pops this async off its dlist can reach the reque...
Definition listen.h:71
void const * app_instance
Definition listen.h:39
fr_app_t const * app
Definition listen.h:38
void const * app_io_instance
I/O path configuration context.
Definition listen.h:33
CONF_SECTION * server_cs
CONF_SECTION of the server.
Definition listen.h:42
fr_dict_t const * dict
dictionary for this listener
Definition listen.h:30
fr_app_io_t const * app_io
I/O path functions.
Definition listen.h:32
Minimal data structure to use the new code.
Definition listen.h:63
unsigned int uint32_t
long int ssize_t
fr_message_set_t * fr_message_set_create(TALLOC_CTX *ctx, int num_messages, size_t message_size, size_t ring_buffer_size, bool unlimited_size)
Create a message set.
Definition message.c:127
int fr_message_done(fr_message_t *m)
Mark a message as done.
Definition message.c:195
fr_message_t * fr_message_and_data_commit(fr_message_set_t *ms, fr_message_t *m, size_t total_size)
Commit a previously reserved message, allocating exactly total_size bytes of packet data.
Definition message.c:1051
void fr_message_set_gc(fr_message_set_t *ms)
Garbage collect the message set.
Definition message.c:1321
fr_message_t * fr_message_and_data_reserve(fr_message_set_t *ms, size_t reserve_size)
Reserve a message.
Definition message.c:973
A Message set, composed of message headers and ring buffer data.
Definition message.c:94
size_t rb_size
cache-aligned size in the ring buffer
Definition message.h:51
fr_time_t when
when this message was sent
Definition message.h:47
uint8_t * data
pointer to the data in the ring buffer
Definition message.h:49
size_t data_size
size of the data in the ring buffer
Definition message.h:50
fr_cmp_ret_t
Result of an ordering comparison.
Definition misc.h:50
static const conf_parser_t config[]
Definition base.c:162
#define fr_assert(_expr)
Definition rad_assert.h:37
#define REDEBUG(fmt,...)
#define RDEBUG(fmt,...)
#define DEBUG2(fmt,...)
static void send_reply(int sockfd, fr_channel_data_t *reply)
int fr_rb_remove(void **removed, fr_rb_tree_t *tree, void const *data)
Remove an entry from the tree, without freeing the data.
Definition rb.c:718
int fr_rb_find(void **found, fr_rb_tree_t const *tree, void const *data)
Find an element in the tree, returning the data, not the node.
Definition rb.c:586
int fr_rb_delete(fr_rb_tree_t *tree, void const *data)
Remove node and free data (if a free function was specified)
Definition rb.c:767
int fr_rb_insert(fr_rb_tree_t *tree, void const *data)
Insert data into a tree.
Definition rb.c:637
#define fr_rb_inline_talloc_alloc(_ctx, _type, _field, _data_cmp, _data_free)
Allocs a red black that verifies elements are of a specific talloc type.
Definition rb.h:244
static bool fr_rb_node_inline_in_tree(fr_rb_node_t const *node)
Check to see if an item is in a tree by examining its inline fr_rb_node_t.
Definition rb.h:312
The main red black tree structure.
Definition rb.h:71
rlm_rcode_t
Return codes indicating the result of the module call.
Definition rcode.h:44
@ RLM_MODULE_TIMEOUT
Module (or section) timed out.
Definition rcode.h:56
int request_slab_deinit(request_t *request)
Callback for slabs to deinitialise the request.
Definition request.c:385
int request_detach(request_t *child)
Unlink a subrequest from its parent.
Definition request.c:544
#define REQUEST_VERIFY(_x)
Definition request.h:310
#define request_is_detached(_x)
Definition request.h:187
#define request_is_external(_x)
Definition request.h:185
#define request_is_internal(_x)
Definition request.h:186
@ REQUEST_TYPE_EXTERNAL
A request received on the wire.
Definition request.h:179
#define request_is_detachable(_x)
Definition request.h:188
#define REQUEST_POOL_NUM_OBJECTS
Definition request.h:68
#define request_init(_ctx, _type, _args)
Definition request.h:322
#define REQUEST_POOL_SIZE
Definition request.h:81
Optional arguments for initialising requests.
Definition request.h:288
fr_ring_buffer_t * fr_ring_buffer_create(TALLOC_CTX *ctx, size_t size)
Create a ring buffer.
Definition ring_buffer.c:64
static char const * name
@ FR_SIGNAL_DUP
A duplicate request was received.
Definition signal.h:44
@ FR_SIGNAL_CANCEL
Request has been cancelled.
Definition signal.h:40
#define FR_SLAB_FUNCS(_name, _type)
Define type specific wrapper functions for slabs and slab elements.
Definition slab.h:124
#define FR_SLAB_TYPES(_name, _type)
Define type specific wrapper structs for slabs and slab elements.
Definition slab.h:75
unsigned int num_children
How many child allocations are expected off each element.
Definition slab.h:48
size_t child_pool_size
Size of pool space to be allocated to each element.
Definition slab.h:49
@ memory_order_seq_cst
Definition stdatomic.h:132
#define atomic_fetch_add_explicit(object, operand, order)
Definition stdatomic.h:302
#define _Atomic(T)
Definition stdatomic.h:77
Definition log.h:93
#define fr_table_str_by_value(_table, _number, _def)
Convert an integer to a string.
Definition table.h:804
static int talloc_const_free(void const *ptr)
Free const'd memory.
Definition talloc.h:288
#define talloc_strdup(_ctx, _str)
Definition talloc.h:149
Definition testlib.h:54
void fr_time_elapsed_update(fr_time_elapsed_t *elapsed, fr_time_t start, fr_time_t end)
Definition time.c:570
void fr_time_elapsed_fprint(FILE *fp, fr_time_elapsed_t const *elapsed, char const *prefix, int tab_offset)
Definition time.c:615
static fr_time_delta_t fr_time_delta_from_msec(int64_t msec)
Definition time.h:575
static int64_t fr_time_delta_unwrap(fr_time_delta_t time)
Definition time.h:154
#define fr_time_delta_lt(_a, _b)
Definition time.h:285
static fr_time_delta_t fr_time_delta_from_sec(int64_t sec)
Definition time.h:590
#define fr_time_delta_wrap(_time)
Definition time.h:152
#define fr_time_delta_ispos(_a)
Definition time.h:290
#define fr_time_eq(_a, _b)
Definition time.h:241
#define NSEC
Definition time.h:379
#define fr_time_add(_a, _b)
Add a time/time delta together.
Definition time.h:196
#define fr_time_sub(_a, _b)
Subtract one time from another.
Definition time.h:229
static fr_time_delta_t fr_time_delta_div(fr_time_delta_t a, fr_time_delta_t b)
Definition time.h:267
static int8_t fr_time_cmp(fr_time_t a, fr_time_t b)
Compare two fr_time_t values.
Definition time.h:916
A time delta, a difference in time measured in nanoseconds.
Definition time.h:80
"server local" time.
Definition time.h:69
@ FR_TIME_TRACKING_YIELDED
We're currently tracking time in the yielded state.
static void fr_time_tracking_yield(fr_time_tracking_t *tt, fr_time_t now)
Transition to the yielded state, recording the time we just spent running.
static void fr_time_tracking_end(fr_time_delta_t *predicted, fr_time_tracking_t *tt, fr_time_t now)
End time tracking for this entity.
fr_time_delta_t waiting_total
total time spent waiting
fr_time_delta_t running_total
total time spent running
static void fr_time_tracking_start(fr_time_tracking_t *parent, fr_time_tracking_t *tt, fr_time_t now)
Start time tracking for a tracked entity.
static void fr_time_tracking_resume(fr_time_tracking_t *tt, fr_time_t now)
Track that a request resumed.
static void fr_time_tracking_debug(fr_time_tracking_t *tt, FILE *fp)
Print debug information about the time tracking structure.
uint64_t fr_timer_list_num_events(fr_timer_list_t *tl)
Return number of pending events.
Definition timer.c:1154
fr_timer_list_t * fr_timer_list_ordered_alloc(TALLOC_CTX *ctx, fr_timer_list_t *parent)
Allocate a new sorted event timer list.
Definition timer.c:1296
int fr_timer_list_force_run(fr_timer_list_t *tl)
Forcibly run all events in an event loop.
Definition timer.c:922
An event timer list.
Definition timer.c:49
#define fr_timer_in(...)
Definition timer.h:87
static bool fr_timer_armed(fr_timer_t *ev)
Definition timer.h:120
static fr_event_list_t * el
static unsigned count
Definition unittest.c:47
void fr_perror(char const *fmt,...)
Print the current error to stderr with a prefix.
Definition strerror.c:737
#define fr_strerror_const_push(_msg)
Definition strerror.h:227
#define fr_strerror_const(_msg)
Definition strerror.h:223
static fr_slen_t data
Definition value.h:1340
static void worker_channel_callback(void const *data, size_t data_size, fr_time_t now, void *uctx)
Handle a control plane message sent to the worker via a channel.
Definition worker.c:244
fr_heap_t * runnable
current runnable requests which we've spent time processing
Definition worker.c:109
static void worker_request_time_tracking_end(fr_worker_t *worker, request_t *request, fr_time_t now)
Definition worker.c:597
static void _worker_request_yield(request_t *request, UNUSED void *uctx)
Interpreter yielded request.
Definition worker.c:1296
fr_event_list_t * el
our event list
Definition worker.c:105
int fr_worker_pre_event(UNUSED fr_time_t now, UNUSED fr_time_delta_t wake, void *uctx)
Pre-event handler.
Definition worker.c:1604
static void worker_send_reply(fr_worker_t *worker, request_t *request, bool do_not_respond, fr_time_t now)
Send a response packet to the network side.
Definition worker.c:614
fr_channel_t * fr_worker_channel_create(fr_worker_t *worker, TALLOC_CTX *ctx, fr_control_t *master)
Create a channel to the worker.
Definition worker.c:1664
fr_rb_tree_t * listeners
so we can cancel requests when a listener goes away
Definition worker.c:116
static void worker_run_request(fr_worker_t *worker, fr_time_t start)
Run a request.
Definition worker.c:1345
static void worker_exit(fr_worker_t *worker)
Definition worker.c:222
#define WORKER_VERIFY
Definition worker.c:67
bool was_sleeping
used to suppress multiple sleep signals in a row
Definition worker.c:128
static int cmd_stats_worker(FILE *fp, UNUSED FILE *fp_err, void *ctx, fr_cmd_info_t const *info)
Definition worker.c:1756
static void _worker_request_runnable(request_t *request, void *uctx)
Request is now runnable.
Definition worker.c:1285
static char * itoa_internal(TALLOC_CTX *ctx, uint64_t number)
Definition worker.c:743
fr_worker_t * fr_worker_alloc(TALLOC_CTX *ctx, fr_event_list_t *el, char const *name, fr_log_t const *logger, fr_log_lvl_t lvl, fr_worker_config_t *config)
Create a worker.
Definition worker.c:1394
fr_worker_channel_t * channel
list of channels
Definition worker.c:131
char const * name
name of this worker
Definition worker.c:91
uint64_t num_active
number of active requests
Definition worker.c:123
fr_cmd_table_t cmd_worker_table[]
Definition worker.c:1794
static int worker_request_time_tracking_start(fr_worker_t *worker, request_t *request, fr_time_t now)
Start time tracking for a request, and mark it as runnable.
Definition worker.c:567
int fr_worker_stats(fr_worker_t const *worker, int num, uint64_t *stats)
Definition worker.c:1739
static int _worker_request_deinit(request_t *request, UNUSED void *uctx)
Definition worker.c:797
static void _worker_request_done_detached(request_t *request, UNUSED rlm_rcode_t rcode, UNUSED void *uctx)
Detached request (i.e.
Definition worker.c:1220
static void _worker_request_resume(request_t *request, UNUSED void *uctx)
Interpreter is starting to work on request again.
Definition worker.c:1305
static fr_cmp_ret_t worker_dedup_cmp(void const *one, void const *two)
Track a request_t in the "dedup" tree.
Definition worker.c:1029
fr_rb_tree_t * dedup
de-dup tree
Definition worker.c:114
fr_atomic_queue_t * aq_control
atomic queue for control messages sent to me
Definition worker.c:101
static void worker_nak(fr_worker_t *worker, fr_channel_data_t *cd, fr_time_t now)
Send a NAK to the network thread.
Definition worker.c:435
static void worker_request_name_number(request_t *request)
Definition worker.c:784
static void _worker_request_timeout(UNUSED fr_timer_list_t *tl, UNUSED fr_time_t when, void *uctx)
Enforce max_request_time.
Definition worker.c:547
fr_log_t const * log
log destination
Definition worker.c:98
fr_io_stats_t stats
input / output stats
Definition worker.c:118
#define CHECK_CONFIG(_x, _min, _max)
static void _worker_request_detach(request_t *request, void *uctx)
Make us responsible for running the request.
Definition worker.c:1253
static int _fr_worker_rb_free(void *arg)
Definition worker.c:161
fr_time_tracking_t tracking
how much time the worker has spent doing things.
Definition worker.c:126
static void _worker_request_done_external(request_t *request, UNUSED rlm_rcode_t rcode, void *uctx)
External request is now complete.
Definition worker.c:1132
void fr_worker_destroy(fr_worker_t *worker)
Destroy a worker.
Definition worker.c:1052
static fr_cmp_ret_t worker_runnable_cmp(void const *one, void const *two)
Track a request_t in the "runnable" heap.
Definition worker.c:1003
uint64_t num_naks
number of messages which were nak'd
Definition worker.c:122
static void worker_request_init(fr_worker_t *worker, request_t *request, fr_time_t now)
Initialize various request fields needed by the worker.
Definition worker.c:766
fr_worker_config_t config
external configuration
Definition worker.c:92
fr_listen_t const * listener
incoming packets
Definition worker.c:137
unlang_interpret_t * intp
Worker's local interpreter.
Definition worker.c:94
static int fr_worker_listen_cancel_self(fr_worker_t *worker, fr_listen_t const *li)
Definition worker.c:386
static void worker_stop_request(request_t *request)
Signal the unlang interpreter that it needs to stop running the request.
Definition worker.c:528
static void _worker_request_prioritise(request_t *request, void *uctx)
Update a request's priority.
Definition worker.c:1322
bool exiting
are we exiting?
Definition worker.c:129
fr_log_lvl_t lvl
log level
Definition worker.c:99
static void worker_requests_cancel(fr_worker_channel_t *ch)
Definition worker.c:213
int num_channels
actual number of channels
Definition worker.c:107
fr_time_delta_t max_request_time
maximum time a request can be processed
Definition worker.c:112
static void worker_recv_request(fr_channel_t *ch, fr_channel_data_t *cd, void *uctx)
Callback which handles a message being received on the worker side.
Definition worker.c:202
static void worker_request_bootstrap(fr_worker_channel_t *wc, fr_channel_data_t *cd, fr_time_t now)
Definition worker.c:802
fr_time_elapsed_t cpu_time
histogram of total CPU time per request
Definition worker.c:119
fr_rb_node_t node
in tree of listeners
Definition worker.c:139
int fr_worker_listen_cancel(fr_worker_t *worker, fr_listen_t const *li)
Definition worker.c:1690
void fr_worker_post_event(UNUSED fr_event_list_t *el, UNUSED fr_time_t now, void *uctx)
Post-event handler.
Definition worker.c:1625
fr_dlist_head_t dlist
of requests associated with this listener.
Definition worker.c:145
void fr_worker(fr_worker_t *worker)
The main loop and entry point of the stand-alone worker thread.
Definition worker.c:1546
request_slab_list_t * slab
slab allocator for request_t
Definition worker.c:133
static uint32_t worker_num_requests(fr_worker_t *worker)
Definition worker.c:792
fr_time_delta_t predicted
How long we predict a request will take to execute.
Definition worker.c:125
pthread_t thread_id
my thread ID
Definition worker.c:96
fr_time_elapsed_t wall_clock
histogram of wall clock time per request
Definition worker.c:120
static bool is_worker_thread(fr_worker_t const *worker)
Definition worker.c:188
fr_worker_channel_t
Definition worker.c:85
fr_timer_list_t * timeout
Track when requests timeout using a dlist.
Definition worker.c:111
static fr_ring_buffer_t * fr_worker_rb_init(void)
Initialise thread local storage.
Definition worker.c:170
fr_control_t * control
the control plane
Definition worker.c:103
static bool _worker_request_scheduled(request_t const *request, UNUSED void *uctx)
Check if a request is scheduled.
Definition worker.c:1314
static void _worker_request_done_internal(request_t *request, UNUSED rlm_rcode_t rcode, void *uctx)
Internal request (i.e.
Definition worker.c:1204
static void worker_listen_cancel_callback(void const *data, NDEBUG_UNUSED size_t data_size, UNUSED fr_time_t now, void *uctx)
A socket is going away, so clean up any requests which use this socket.
Definition worker.c:413
void fr_worker_debug(fr_worker_t *worker, FILE *fp)
Print debug information about the worker structure.
Definition worker.c:1637
static void _worker_request_internal_init(request_t *request, void *uctx)
Internal request (i.e.
Definition worker.c:1111
static fr_cmp_ret_t worker_listener_cmp(void const *one, void const *two)
Definition worker.c:149
#define CHECK_CONFIG_TIME_DELTA(_x, _min, _max)
A worker which takes packets from a master, and processes them.
Definition worker.c:90
int message_set_size
default start number of messages
Definition worker.h:73
#define FR_CONTROL_ID_LISTEN_DEAD
Definition worker.h:40
int max_requests
max requests this worker will handle
Definition worker.h:69
int max_channels
maximum number of channels
Definition worker.h:71
fr_slab_config_t reuse
slab allocator configuration
Definition worker.h:78
int ring_buffer_size
default start size for the ring buffers
Definition worker.h:74
fr_time_delta_t max_request_time
maximum time a request can be processed
Definition worker.h:76