The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
network.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: b005475a2209b501780c782afd219e6bad6f0883 $
19 *
20 * @brief Receiver of socket data, which sends messages to the workers.
21 * @file io/network.c
22 *
23 * @copyright 2016 Alan DeKok (aland@freeradius.org)
24 */
25RCSID("$Id: b005475a2209b501780c782afd219e6bad6f0883 $")
26
27#define LOG_PREFIX nr->name
28
29#define LOG_DST nr->log
30
31#include <freeradius-devel/util/event.h>
32#include <freeradius-devel/util/rand.h>
33#include <freeradius-devel/util/rb.h>
34#include <freeradius-devel/util/syserror.h>
35#include <freeradius-devel/util/atexit.h>
36
37#include <freeradius-devel/io/channel.h>
38#include <freeradius-devel/io/listen.h>
39#include <freeradius-devel/io/network.h>
40#include <freeradius-devel/io/queue.h>
41
42#define MAX_WORKERS 64
43
45
52
53/** Associate a worker thread with a network thread
54 *
55 */
56typedef struct {
57 fr_heap_index_t heap_id; //!< workers are in a heap
58 fr_time_delta_t cpu_time; //!< how much CPU time this worker has spent
59 fr_time_delta_t predicted; //!< predicted processing time for one packet
60
61 bool blocked; //!< is this worker blocked?
62
63 fr_channel_t *channel; //!< channel to the worker
64 fr_message_set_t *reply_ms; //!< message set the worker will use to send reply data.
65 fr_worker_t *worker; //!< worker pointer
68
69typedef struct {
70 fr_rb_node_t listen_node; //!< rbtree node for looking up by listener.
71 fr_rb_node_t num_node; //!< rbtree node for looking up by number.
72
73 fr_network_t *nr; //!< O(N) issues in talloc
74 int number; //!< unique ID
75 fr_heap_index_t heap_id; //!< for the sockets_by_num heap
76
77 fr_event_filter_t filter; //!< what type of filter it is
78 fr_event_fd_t *ef; //!< the I/O event, if it was inserted. Parented off
79 ///< this socket, so it cannot outlive it.
80
81 bool dead; //!< is it dead?
82 bool blocked; //!< is it blocked?
83 bool closed; //!< the descriptor has been closed.
84
85 unsigned int outstanding; //!< number of outstanding packets sent to the worker
86 fr_listen_t *listen; //!< I/O ctx and functions.
87
88 fr_message_set_t *ms; //!< message buffers for this socket.
89 fr_channel_data_t *cd; //!< cached in case of allocation & read error
90 size_t leftover; //!< leftover data from a previous read
91 size_t written; //!< however much we did in a partial write
92
93 fr_channel_data_t *pending; //!< the currently pending partial packet
94 fr_heap_t *waiting; //!< packets waiting to be written
97
98/*
99 * We have an array of workers, so we can index the workers in
100 * O(1) time. remove the heap of "workers ordered by CPU time"
101 * when we send a packet to a worker, just update the predicted
102 * CPU time in place. when we receive a reply from a worker,
103 * just update the predicted CPU time in place.
104 *
105 * when we need to choose a worker, pick 2 at random, and then
106 * choose the one with the lowe cpu time. For background, see
107 * "Power of Two-Choices" and
108 * https://www.eecs.harvard.edu/~michaelm/postscripts/mythesis.pdf
109 * https://www.eecs.harvard.edu/~michaelm/postscripts/tpds2001.pdf
110 */
112 char const *name; //!< Network ID for logging.
113
114 pthread_t thread_id; //!< for self
115
116 bool suspended; //!< whether or not we're suspended.
117
118 fr_log_t const *log; //!< log destination
119 fr_log_lvl_t lvl; //!< debug log level
120
121 fr_atomic_queue_t *aq_control; //!< atomic queue for control messages sent to me
122
123 fr_control_t *control; //!< the control plane
124
125 fr_ring_buffer_t *rb; //!< ring buffer for my control-plane messages
126
127 fr_event_list_t *el; //!< our event list
128
129 fr_heap_t *replies; //!< replies from the worker, ordered by priority / origin time
130
132
133 fr_rb_tree_t *sockets; //!< list of sockets we're managing, ordered by the listener
134 fr_rb_tree_t *sockets_by_num; //!< ordered by number;
135
136 int num_workers; //!< number of active workers
137 int num_blocked; //!< number of blocked workers
138 int num_pending_workers; //!< number of workers we're waiting to start.
139 int max_workers; //!< maximum number of allowed workers
140 int num_sockets; //!< actually a counter...
141
142 int signal_pipe[2]; //!< Pipe for signalling the worker in an orderly way.
143 ///< This is more deterministic than using async signals.
144
145 bool exiting; //!< are we exiting?
146
147 fr_network_config_t config; //!< configuration
149};
150
151static void fr_network_post_event(fr_event_list_t *el, fr_time_t now, void *uctx);
152static int fr_network_pre_event(fr_time_t now, fr_time_delta_t wake, void *uctx);
155static void fr_network_read(UNUSED fr_event_list_t *el, int sockfd, UNUSED int flags, void *ctx);
156
157static fr_cmp_ret_t reply_cmp(void const *one, void const *two)
158{
159 fr_channel_data_t const *a = one, *b = two;
160 int ret;
161
162 ret = CMP(a->priority, b->priority);
163 if (ret != 0) return ret;
164
165 return fr_time_cmp(a->m.when, b->m.when);
166}
167
168static fr_cmp_ret_t waiting_cmp(void const *one, void const *two)
169{
170 fr_channel_data_t const *a = one, *b = two;
171 int ret;
172
173 ret = CMP(a->priority, b->priority);
174 if (ret != 0) return ret;
175
176 return fr_time_cmp(a->reply.request_time, b->reply.request_time);
177}
178
179static fr_cmp_ret_t socket_listen_cmp(void const *one, void const *two)
180{
181 fr_network_socket_t const *a = one, *b = two;
182
183 return CMP(a->listen, b->listen);
184}
185
186static fr_cmp_ret_t socket_num_cmp(void const *one, void const *two)
187{
188 fr_network_socket_t const *a = one, *b = two;
189
190 return CMP(a->number, b->number);
191}
192
193/*
194 * Explicitly cleanup the memory allocated to the ring buffer,
195 * just in case valgrind complains about it.
196 */
197static int _fr_network_rb_free(void *arg)
198{
199 return talloc_free(arg);
200}
201
202/** Initialise thread local storage
203 *
204 * @return fr_ring_buffer_t for messages
205 */
207{
209
210 rb = fr_network_rb;
211 if (rb) return rb;
212
214 if (!rb) {
215 fr_perror("Failed allocating memory for network ring buffer");
216 return NULL;
217 }
218
220
221 return rb;
222}
223
224static inline bool is_network_thread(fr_network_t const *nr)
225{
226 return (pthread_equal(pthread_self(), nr->thread_id) != 0);
227}
228
230
231/** Add a fr_listen_t to a network
232 *
233 * @param nr the network
234 * @param li the listener
235 */
237{
239
240 /*
241 * Associate the protocol dictionary with the listener, so that the decode functions can check /
242 * use it.
243 *
244 * A virtual server may start off with a "dictionary" block, and therefore define a local
245 * dictionary. So the "root" dictionary of a virtual server may not be a protocol dict.
246 */
247 fr_assert(li->server_cs != NULL);
249
250 fr_assert(li->dict != NULL);
251
252 /*
253 * Skip a bunch of work if we're already in the network thread.
254 */
255 if (is_network_thread(nr) && !li->needs_full_setup) {
256 return fr_network_listen_add_self(nr, li);
257 }
258
259 rb = fr_network_rb_init();
260 if (!rb) return -1;
261
262 return fr_control_message_send(nr->control, rb, FR_CONTROL_ID_LISTEN, &li, sizeof(li));
263}
264
265
266/** Delete a socket from a network. MUST be called only by the listener itself!.
267 *
268 * @param nr the network
269 * @param li the listener
270 */
272{
274
276
277 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
278 if (!s) return -1;
279
281
282 return 0;
283}
284
285/** Add a "watch directory" call to a network
286 *
287 * @param nr the network
288 * @param li the listener
289 */
291{
293
294 rb = fr_network_rb_init();
295 if (!rb) return -1;
296
297 return fr_control_message_send(nr->control, rb, FR_CONTROL_ID_DIRECTORY, &li, sizeof(li));
298}
299
300/** Add a worker to a network in a different thread
301 *
302 * @param nr the network
303 * @param worker the worker
304 */
306{
308
309 rb = fr_network_rb_init();
310 if (!rb) return -1;
311
312 (void) talloc_get_type_abort(nr, fr_network_t);
313 (void) talloc_get_type_abort(worker, fr_worker_t);
314
315 return fr_control_message_send(nr->control, rb, FR_CONTROL_ID_WORKER, &worker, sizeof(worker));
316}
317
318static void fr_network_worker_started_callback(void const *data, size_t data_size, fr_time_t now, void *uctx);
319
320/** Add a worker to a network in the same thread
321 *
322 * @param nr the network
323 * @param worker the worker
324 */
326{
327 fr_network_worker_started_callback(&worker, sizeof(worker), fr_time_wrap(0), nr);
328}
329
330
331/** Signal the network to read from a listener
332 *
333 * @param nr the network
334 * @param li the listener to read from
335 */
337{
339
340 (void) talloc_get_type_abort(nr, fr_network_t);
342
343 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
344 if (!s) return;
345
346 /*
347 * Go read the socket.
348 */
349 fr_network_read(nr->el, s->listen->fd, 0, s);
350}
351
352
353/** Inject a packet for a listener to write
354 *
355 * @param nr the network
356 * @param li the listener where the packet is being injected
357 * @param packet the packet to be written
358 * @param packet_len the length of the packet
359 * @param packet_ctx The packet context to write
360 * @param request_time when the packet was received.
361 */
362void fr_network_listen_write(fr_network_t *nr, fr_listen_t *li, uint8_t const *packet, size_t packet_len,
363 void *packet_ctx, fr_time_t request_time)
364{
365 fr_message_t *lm;
367
368 cd = (fr_channel_data_t) {
369 .m = (fr_message_t) {
371 .data_size = packet_len,
372 .when = request_time,
373 },
374
375 .channel = {
376 .heap_id = FR_HEAP_INDEX_INVALID,
377 },
378
379 .listen = li,
380 .priority = PRIORITY_NOW,
381 .reply.request_time = request_time,
382 };
383
384 memcpy(&cd.m.data, &packet, sizeof(packet)); /* const issues */
385 memcpy(&cd.packet_ctx, &packet_ctx, sizeof(packet_ctx)); /* const issues */
386
387 /*
388 * Localize the message and insert it into the heap of pending messages.
389 */
390 lm = fr_message_localize(nr, &cd.m, sizeof(cd));
391 if (!lm) return;
392
393 if (fr_heap_insert(&nr->replies, lm) < 0) {
394 fr_message_done(lm);
395 }
396}
397
398
399/** Inject a packet for a listener to read
400 *
401 * @param nr the network
402 * @param li the listener where the packet is being injected
403 * @param packet the packet to be injected
404 * @param packet_len the length of the packet
405 * @param recv_time when the packet was received.
406 * @return
407 * - <0 on error
408 * - 0 on success
409 */
410int fr_network_listen_inject(fr_network_t *nr, fr_listen_t *li, uint8_t const *packet, size_t packet_len, fr_time_t recv_time)
411{
412 int rcode;
414 fr_network_inject_t my_inject;
415
416 /*
417 * Can't inject to injection-less destinations.
418 */
419 if (!li->app_io->inject) {
420 fr_strerror_const("Listener cannot accept injected packet");
421 return -1;
422 }
423
424 /*
425 * Avoid a bounce through the event loop if we're being called from the network thread.
426 */
427 if (is_network_thread(nr)) {
429
430 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
431 if (!s) {
432 fr_strerror_const("Listener was not found for injected packet");
433 return -1;
434 }
435
436 /*
437 * Inject the packet. The master.c mod_read() routine will then take care of avoiding
438 * IO, and instead return the packet to the network side.
439 */
440 if (li->app_io->inject(li, packet, packet_len, recv_time) == 0) {
441 (void) fr_network_read(nr->el, li->fd, 0, s);
442 }
443
444 return 0;
445 }
446
447 rb = fr_network_rb_init();
448 if (!rb) return -1;
449
450 my_inject.listen = li;
451 MEM(my_inject.packet = talloc_memdup(NULL, packet, packet_len));
452 my_inject.packet_len = packet_len;
453 my_inject.recv_time = recv_time;
454
455 rcode = fr_control_message_send(nr->control, rb, FR_CONTROL_ID_INJECT, &my_inject, sizeof(my_inject));
456 if (rcode < 0) talloc_free(my_inject.packet);
457
458 return rcode;
459}
460
463 { 0 }
464};
465
468 { 0 }
469};
470
473 { 0 }
474};
475
478 { 0 }
479};
480
482{
485
486 if (nr->suspended) return;
487
488 for (s = fr_rb_iter_init_inorder(nr->sockets, &iter);
489 s != NULL;
490 s = fr_rb_iter_next_inorder(nr->sockets, &iter)) {
492 }
493 nr->suspended = true;
494}
495
497{
500
501 if (!nr->suspended) return;
502
503 for (s = fr_rb_iter_init_inorder(nr->sockets, &iter);
504 s != NULL;
505 s = fr_rb_iter_next_inorder(nr->sockets, &iter)) {
507 }
508 nr->suspended = false;
509}
510
511#define IALPHA (8)
512#define RTT(_old, _new) fr_time_delta_wrap((fr_time_delta_unwrap(_new) + (fr_time_delta_unwrap(_old) * (IALPHA - 1))) / IALPHA)
513
514/** Callback which handles a message being received on the network side.
515 *
516 * @param[in] ch the channel that the message is on.
517 * @param[in] cd the message (if any) to start with
518 * @param[in] uctx the network
519 */
520static void fr_network_recv_reply(fr_channel_t *ch, fr_channel_data_t *cd, void *uctx)
521{
522 fr_network_t *nr = uctx;
523 fr_network_worker_t *worker;
524
525 cd->channel.ch = ch;
526
527 /*
528 * Update stats for the worker.
529 */
531 worker->stats.out++;
532 worker->cpu_time = cd->reply.cpu_time;
533 if (!fr_time_delta_ispos(worker->predicted)) {
534 worker->predicted = cd->reply.processing_time;
535 } else {
536 worker->predicted = RTT(worker->predicted, cd->reply.processing_time);
537 }
538
539 /*
540 * Unblock the worker.
541 */
542 if (worker->blocked) {
543 worker->blocked = false;
544 nr->num_blocked--;
546 }
547
548 /*
549 * Ensure that heap insert works.
550 */
551 cd->channel.heap_id = FR_HEAP_INDEX_INVALID;
552 if (fr_heap_insert(&nr->replies, cd) < 0) {
553 fr_message_done(&cd->m);
554 fr_assert(0 == 1);
555 }
556}
557
558/** Handle a network control message callback for a channel
559 *
560 * This is called from the event loop when we get a notification
561 * from the event signalling pipe.
562 *
563 * @param[in] data the message
564 * @param[in] data_size size of the data
565 * @param[in] now the current time
566 * @param[in] uctx the network
567 */
568static void fr_network_channel_callback(void const *data, size_t data_size, fr_time_t now, void *uctx)
569{
571 fr_channel_t *ch;
572 fr_network_t *nr = uctx;
573
574 ce = fr_channel_service_message(now, &ch, NULL, data, data_size);
575 DEBUG3("Channel %s",
576 fr_table_str_by_value(channel_signals, ce, "<INVALID>"));
577 switch (ce) {
578 case FR_CHANNEL_ERROR:
579 return;
580
581 case FR_CHANNEL_EMPTY:
582 return;
583
584 case FR_CHANNEL_NOOP:
585 break;
586
588 fr_assert(ch != NULL);
589 while (fr_channel_recv_reply(ch));
590 break;
591
593 fr_assert(0 == 1);
594 break;
595
596 case FR_CHANNEL_OPEN:
597 fr_assert(0 == 1);
598 break;
599
600 case FR_CHANNEL_CLOSE:
601 {
602 fr_network_worker_t *w = talloc_get_type_abort(fr_channel_requestor_uctx_get(ch),
604 int i;
605
606 /*
607 * Remove this worker from the array
608 */
609 DEBUG3("Worker acked our close request");
610 for (i = 0; i < nr->num_workers; i++) {
611 if (nr->workers[i] == w) {
612 if (i == (nr->num_workers - 1)) break;
613
614 /*
615 * Close the hole...
616 */
617 memmove(&nr->workers[i], &nr->workers[i + 1],
618 (uint8_t *) &nr->workers[nr->num_workers] - (uint8_t *) &nr->workers[i + 1]);
619 break;
620 }
621 }
622 nr->num_workers--;
623 nr->workers[nr->num_workers] = NULL; /* over-write now unused pointer */
624 }
625 break;
626 }
627}
628
629#define OUTSTANDING(_x) ((_x)->stats.in - (_x)->stats.out)
630
631/** Send a message on the "best" channel.
632 *
633 * @param nr the network
634 * @param cd the message we've received
635 */
637{
638 fr_network_worker_t *worker;
639
640 (void) talloc_get_type_abort(nr, fr_network_t);
641
642 /*
643 * The workers have been signalled to close and are tearing their
644 * channels down. Anything queued now is stranded: the worker
645 * discards it without a reply, so our outstanding count for the
646 * socket never comes back down and the socket can never be freed.
647 *
648 * The listeners are closed before the close is signalled, so a
649 * socket read cannot get here. Assert, so that whatever did shows
650 * itself with a backtrace, and drop the packet in release builds.
651 */
652 if (!fr_cond_assert_msg(!nr->exiting,
653 "Sending packet to worker after signalling the channel close")) return -1;
654
655 if (!nr->num_workers) {
656 RATE_LIMIT_GLOBAL(ERROR, "Failed sending packet to worker - "
657 "No workers are available");
658 return -1;
659 }
660
661retry:
662 if (nr->num_workers == 1) {
663 worker = nr->workers[0];
664 if (worker->blocked) {
665 RATE_LIMIT_GLOBAL(ERROR, "Failed sending packet to worker - "
666 "In single-threaded mode and worker is blocked");
667 drop:
668 worker->stats.dropped++;
669 return -1;
670 }
671
672 } else if (nr->num_blocked == 0) {
673 int64_t cmp;
674 uint32_t one, two;
675
676 one = fr_rand() % nr->num_workers;
677 do {
678 two = fr_rand() % nr->num_workers;
679 } while (two == one);
680
681 /*
682 * Choose a worker based on minimizing the amount
683 * of future work it's being asked to do.
684 *
685 * If both workers have the same number of
686 * outstanding requests, then choose the worker
687 * which has used the least total CPU time.
688 */
689 cmp = (OUTSTANDING(nr->workers[one]) - OUTSTANDING(nr->workers[two]));
690 if (cmp < 0) {
691 worker = nr->workers[one];
692
693 } else if (cmp > 0) {
694 worker = nr->workers[two];
695
696 } else if (fr_time_delta_lt(nr->workers[one]->cpu_time, nr->workers[two]->cpu_time)) {
697 worker = nr->workers[one];
698
699 } else {
700 worker = nr->workers[two];
701 }
702 } else {
703 int i;
704 uint64_t min_outstanding = UINT64_MAX;
705 fr_network_worker_t *found = NULL;
706
707 /*
708 * Some workers are blocked. Pick the worker
709 * with the least amount of future work to do.
710 */
711 for (i = 0; i < nr->num_workers; i++) {
712 uint64_t outstanding;
713
714 worker = nr->workers[i];
715 if (worker->blocked) continue;
716
717 outstanding = OUTSTANDING(worker);
718 if ((outstanding < min_outstanding) || !found) {
719 found = worker;
720 min_outstanding = outstanding;
721
722 } else if (outstanding == min_outstanding) {
723 /*
724 * Queue lengths are the same.
725 * Choose this worker if it's
726 * less busy than the previous one we found.
727 */
728 if (fr_time_delta_lt(worker->cpu_time, found->cpu_time)) {
729 found = worker;
730 }
731 }
732 }
733
734 if (!found) {
735 RATE_LIMIT_GLOBAL(PERROR, "Failed sending packet to worker - Couldn't find active worker, "
736 "%u/%u workers are blocked", nr->num_blocked, nr->num_workers);
737 return -1;
738 }
739
740 worker = found;
741 }
742
743 (void) talloc_get_type_abort(worker, fr_network_worker_t);
744
745 /*
746 * Too many outstanding packets for this worker. Drop
747 * the request.
748 *
749 * If the worker we've picked has too many outstanding
750 * packets, then we have either only one worker, in which
751 * cae we should drop the packet. Or, we were unable to
752 * find a worker with smaller than max_outstanding
753 * packets. In which case all of the workers are likely
754 * at max_outstanding.
755 *
756 * In both cases, we should just drop the new packet.
757 */
758 fr_assert(worker->stats.in >= worker->stats.out);
759 if (nr->config.max_outstanding &&
760 (OUTSTANDING(worker) >= nr->config.max_outstanding)) {
761 RATE_LIMIT_GLOBAL(PERROR, "max_outstanding reached - dropping packet");
762 goto drop;
763 }
764
765 /*
766 * Send the message to the channel. If we fail, drop the
767 * packet. The only reason for failure is that the
768 * worker isn't servicing it's input queue. When that
769 * happens, we have no idea what to do, and the whole
770 * thing falls over.
771 */
772 if (fr_channel_send_request(worker->channel, cd) < 0) {
773 worker->stats.dropped++;
774 worker->blocked = true;
775 nr->num_blocked++;
776
777 RATE_LIMIT_GLOBAL(PERROR, "Failed sending packet to worker - %u/%u workers are blocked",
778 nr->num_blocked, nr->num_workers);
779
780 if (nr->num_blocked == nr->num_workers) {
782 return -1;
783 }
784 goto retry;
785 }
786
787 worker->stats.in++;
788
789 /*
790 * We're projecting that the worker will use more CPU
791 * time to process this request. The CPU time will be
792 * updated with a more accurate number when we receive a
793 * reply from this channel.
794 */
795 worker->cpu_time = fr_time_delta_add(worker->cpu_time, worker->predicted);
796
797 return 0;
798}
799
800
801/** Send a packet to the worker.
802 *
803 * MUST only be called from the network thread.
804 *
805 * @param nr the network
806 * @param parent the parent listener
807 * @param li the listener that the packet was "read" from. Can be "parent"
808 * @param buffer the packet to send
809 * @param buflen size of the packet to send
810 * @param recv_time of the packet
811 * @param packet_ctx for the packet
812 * @return
813 * - <0 on error
814 * - 0 on success
815 */
817 const uint8_t *buffer, size_t buflen, fr_time_t recv_time, void *packet_ctx)
818{
821
822 (void) talloc_get_type_abort(nr, fr_network_t);
824
825 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
826 if (!s) return -1;
827
829 if (!cd) return -1;
830
831 cd->listen = parent;
833 cd->packet_ctx = packet_ctx;
834 cd->request.recv_time = recv_time;
835 memcpy(cd->m.data, buffer, buflen);
836 cd->m.when = fr_time();
837
838 if (fr_network_send_request(nr, cd) < 0) {
840 fr_message_done(&cd->m);
841 nr->stats.dropped++;
842 s->stats.dropped++;
843 return -1;
844 }
845
846 s->outstanding++;
847 return 0;
848}
849
850/** Get the number of outstanding packets
851 *
852 * @param nr the network
853 * @param li the listener that the packet was "read" from
854 * @return
855 * - <0 on error
856 * - the number of outstanding packets
857*/
860
861 (void) talloc_get_type_abort(nr, fr_network_t);
863
864 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
865 if (!s) return -1;
866
867 return s->outstanding;
868}
869
870/*
871 * Mark it as dead, but DON'T free it until all of the replies
872 * have come in.
873 */
875{
876 int i;
877
878 if (s->dead) return;
879
880 s->dead = true;
881
882 /*
883 * Nothing writes to a dead socket - fr_network_post_event()
884 * discards its replies - so the descriptor goes now. The memory
885 * has to wait for the workers to finish with its messages.
886 */
887 if (network_socket_close(s) < 0) PWARN("Failed removing event for socket %d", s->number);
888
889 for (i = 0; i < nr->max_workers; i++) {
890 if (!nr->workers[i]) continue;
891
892 (void) fr_worker_listen_cancel(nr->workers[i]->worker, s->listen);
893 }
894
895 /*
896 * If there are no outstanding packets, then we can free
897 * it now.
898 */
899 if (!s->outstanding) {
900 talloc_free(s);
901 return;
902 }
903
904 /*
905 * There are still outstanding packets. Leave it in the
906 * socket tree, so that replies from the worker can find
907 * it. When we've received all of the replies, then
908 * fr_network_post_event() will clean up this socket.
909 */
910}
911
912/** Read a packet from the network.
913 *
914 * @param[in] el the event list.
915 * @param[in] sockfd the socket which is ready to read.
916 * @param[in] flags from kevent.
917 * @param[in] ctx the network socket context.
918 */
919static void fr_network_read(UNUSED fr_event_list_t *el, int sockfd, UNUSED int flags, void *ctx)
920{
921 int num_messages = 0;
922 fr_network_socket_t *s = ctx;
923 fr_network_t *nr = s->nr;
924 ssize_t data_size;
925 fr_channel_data_t *cd, *next;
926
927 if (!fr_cond_assert_msg(s->listen->fd == sockfd, "Expected listen->fd (%u) to be equal event fd (%u)",
928 s->listen->fd, sockfd)) return;
929
930 DEBUG3("Reading data from FD %u", sockfd);
931
932 if (!s->cd) {
934 if (!cd) {
935 ERROR("Failed allocating message size %zd! - Closing socket",
938 return;
939 }
940 } else {
941 cd = s->cd;
942 }
943
944 fr_assert(cd->m.data != NULL);
945
946next_message:
947 /*
948 * Poll this socket, but not too often. We have to go
949 * service other sockets, too.
950 */
951 if (num_messages > 16) {
952 s->cd = cd;
953 return;
954 }
955
957
958 /*
959 * Read data from the network.
960 *
961 * Return of 0 means "no data", which is fine for UDP.
962 * For TCP, if an underlying read() on the TCP socket
963 * returns 0, (which signals that the FD is no longer
964 * usable) this function should return -1, so that the
965 * network side knows that it needs to close the
966 * connection.
967 */
968 data_size = s->listen->app_io->read(s->listen, &cd->packet_ctx, &cd->request.recv_time,
969 cd->m.data, cd->m.rb_size, &s->leftover);
970 if (data_size == 0) {
971 /*
972 * Cache the message only when there are leftover
973 * bytes from a partial stream read. The buffer
974 * must be handed back on the next call so the
975 * stream can append to it.
976 *
977 * When app_io->read() calls fr_network_listen_send_packet()
978 * internally, that function calls fr_message_and_data_alloc()
979 * on the same message set. Because the message ring uses
980 * fr_ring_buffer_reserve() (which does not advance write_offset),
981 * the new alloc returns the same ring slot as our reservation.
982 * The memset inside message_reserve zeroes our struct, then the
983 * new message commits into that slot. We can detect this
984 * because cd->m.data_size is non-zero after the alloc commits.
985 *
986 * In the aliased case, cd now refers to the already-committed
987 * message that was dispatched to a worker. We must not reset
988 * or modify it; just clear s->cd so the next call reserves fresh.
989 *
990 * When neither leftover nor aliasing applies, the reservation
991 * is clean and unused; reset it explicitly so its ring-buffer
992 * slot is immediately reusable.
993 */
994 if (s->leftover) {
995 s->cd = cd;
996 } else if (cd->m.data_size != 0) {
997 s->cd = NULL;
998 } else {
1000 s->cd = NULL;
1001 }
1002 return;
1003 }
1004
1005 /*
1006 * Error: close the connection, and remove the fr_listen_t
1007 */
1008 if (data_size < 0) {
1009// fr_log(nr->log, L_DBG_ERR, "error from transport read on socket %d", sockfd);
1011 return;
1012 }
1013 s->cd = NULL;
1014
1015 DEBUG3("Read %zd byte(s) from FD %u", data_size, sockfd);
1016 if (s->listen->read_hexdump) HEXDUMP2(cd->m.data, data_size, "%s read ", s->listen->name);
1017 nr->stats.in++;
1018 s->stats.in++;
1019
1020 /*
1021 * Initialize the rest of the fields of the channel data.
1022 *
1023 * We always use "now" as the time of the message, as the
1024 * packet MAY be a duplicate packet magically resurrected
1025 * from the past. i.e. If the read routines are doing
1026 * dedup, then they notice that the packet is a
1027 * duplicate. In that case, they send over a copy of the
1028 * packet, BUT with the original timestamp. This
1029 * information tells the worker that the packet is a
1030 * duplicate.
1031 */
1032 cd->m.when = fr_time();
1033 cd->listen = s->listen;
1034
1035 /*
1036 * Nothing in the buffer yet. Allocate room for one
1037 * packet.
1038 */
1039 if ((cd->m.data_size == 0) && (!s->leftover)) {
1040 (void) fr_message_and_data_commit(s->ms, &cd->m, data_size);
1041 next = NULL;
1042 } else {
1043 /*
1044 * There are leftover bytes in the buffer, feed
1045 * them to the next round of reading.
1046 */
1047 if (s->leftover) {
1048 next = (fr_channel_data_t *) fr_message_and_data_commit_with_leftover(s->ms, &cd->m, data_size, s->leftover,
1050 if (!next) {
1051 PERROR("Failed reserving partial packet.");
1052 // @todo - probably close the socket...
1053 fr_assert(0 == 1);
1054 }
1055 } else {
1056 (void) fr_message_and_data_commit(s->ms, &cd->m, data_size);
1057 next = NULL;
1058 }
1059 }
1060
1061 /*
1062 * Set the priority. Which incidentally also checks if
1063 * we're allowed to read this particular kind of packet.
1064 *
1065 * That check is because the app_io handlers just read
1066 * packets, and don't really have access to the parent
1067 * "list of allowed packet types". So we have to do the
1068 * work here in a callback.
1069 *
1070 * That should probably be fixed...
1071 */
1072 if (s->listen->app->priority) {
1073 int priority;
1074
1075 priority = s->listen->app->priority(s->listen->app_instance, cd->m.data, data_size);
1076 if (priority <= 0) goto discard;
1077
1078 cd->priority = priority;
1079 }
1080
1081 if (fr_network_send_request(nr, cd) < 0) {
1082 discard:
1083 talloc_free(cd->packet_ctx); /* not sure what else to do here */
1084 fr_message_done(&cd->m);
1085 nr->stats.dropped++;
1086 s->stats.dropped++;
1087
1088 } else {
1089 /*
1090 * One more packet sent to a worker.
1091 */
1092 s->outstanding++;
1093 }
1094
1095 /*
1096 * If there is a next message, go read it from the buffer.
1097 *
1098 * @todo - note that this calls read(), even if the
1099 * app_io has paused the reader. We likely want to be
1100 * able to check that, too. We might just remove this
1101 * "goto"...
1102 */
1103 if (next) {
1104 cd = next;
1105 num_messages++;
1106 goto next_message;
1107 }
1108}
1109
1110int fr_network_sendto_worker(fr_network_t *nr, fr_listen_t *li, void *packet_ctx, uint8_t const *data, size_t data_len, fr_time_t recv_time)
1111{
1114
1115 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
1116 if (!s) return -1;
1117
1118 cd = (fr_channel_data_t *) fr_message_and_data_alloc(s->ms, data_len);
1119 if (!cd) return -1;
1120
1121 s->stats.in++;
1122
1124
1125 cd->m.when = recv_time;
1126 cd->listen = li;
1127 cd->packet_ctx = packet_ctx;
1128
1129 memcpy(cd->m.data, data, data_len);
1130
1131 if (fr_network_send_request(nr, cd) < 0) {
1132 talloc_free(packet_ctx);
1133 fr_message_done(&cd->m);
1134 nr->stats.dropped++;
1135 s->stats.dropped++;
1136 return -1;
1137 }
1138
1139 /*
1140 * One more packet sent to a worker.
1141 */
1142 s->outstanding++;
1143 return 0;
1144}
1145
1146
1147/** Get a notification that a vnode changed
1148 *
1149 * @param[in] el the event list.
1150 * @param[in] sockfd the socket which is ready to read.
1151 * @param[in] fflags from kevent.
1152 * @param[in] ctx the network socket context.
1153 */
1154static void fr_network_vnode_extend(UNUSED fr_event_list_t *el, int sockfd, int fflags, void *ctx)
1155{
1156 fr_network_socket_t *s = ctx;
1157 fr_network_t *nr = s->nr;
1158
1159 if (!fr_cond_assert(s->listen->fd == sockfd)) return;
1160
1161 DEBUG3("network vnode");
1162
1163 /*
1164 * Tell the IO handler that something has happened to the
1165 * file.
1166 */
1167 s->listen->app_io->vnode(s->listen, fflags);
1168}
1169
1170
1171/** Handle errors for a socket.
1172 *
1173 * @param[in] el the event list
1174 * @param[in] sockfd the socket which has a fatal error.
1175 * @param[in] flags returned by kevent.
1176 * @param[in] fd_errno returned by kevent.
1177 * @param[in] ctx the network socket context.
1178 */
1180 int fd_errno, void *ctx)
1181{
1182 fr_network_socket_t *s = ctx;
1183 fr_network_t *nr = s->nr;
1184
1185 if (s->listen->app_io->error) {
1186 s->listen->app_io->error(s->listen);
1187
1188 } else if (flags & EV_EOF) {
1189 DEBUG2("Socket %s closed by peer", s->listen->name);
1190
1191 } else {
1192 ERROR("Socket %s errored - %s", s->listen->name, fr_syserror(fd_errno));
1193 }
1194
1196}
1197
1198
1199/** Write packets to the network.
1200 *
1201 * @param el the event list
1202 * @param sockfd the socket which is ready to write
1203 * @param flags returned by kevent.
1204 * @param ctx the network socket context.
1205 */
1206static void fr_network_write(UNUSED fr_event_list_t *el, UNUSED int sockfd, UNUSED int flags, void *ctx)
1207{
1208 fr_network_socket_t *s = ctx;
1209 fr_listen_t *li = s->listen;
1210 fr_network_t *nr = s->nr;
1212
1213 (void) talloc_get_type_abort(nr, fr_network_t);
1214
1215 /*
1216 * Start with the currently pending message, and then
1217 * work through the priority heap.
1218 */
1219 if (s->pending) {
1220 cd = s->pending;
1221 s->pending = NULL;
1222
1223 } else {
1224 fr_heap_pop((void **)&cd, &s->waiting);
1225 }
1226
1227 while (cd != NULL) {
1228 int rcode;
1229
1230 fr_assert(li == cd->listen);
1231 if (li->write_hexdump) HEXDUMP2(cd->m.data, cd->m.data_size, "%s writing ", li->name);
1232 rcode = li->app_io->write(li, cd->packet_ctx,
1233 cd->reply.request_time,
1234 cd->m.data, cd->m.data_size, s->written);
1235
1236 /*
1237 * Write of 0 bytes means an OS bug, and we just discard this packet.
1238 */
1239 if (rcode == 0) {
1240 RATE_LIMIT_GLOBAL(ERROR, "Discarding packet due to write returning zero for socket %s",
1241 s->listen->name);
1242 goto discard;
1243 }
1244
1245 /*
1246 * Or we have a write error.
1247 */
1248 if (rcode < 0) {
1249 /*
1250 * Stop processing the heap, and set the
1251 * pending message to the current one.
1252 */
1253 if (errno == EWOULDBLOCK) {
1254 save_pending:
1255 fr_assert(!s->pending);
1256
1257 if (cd->m.status != FR_MESSAGE_LOCALIZED) {
1258 fr_message_t *lm;
1259
1260 lm = fr_message_localize(s, &cd->m, sizeof(*cd));
1261 if (!lm) {
1262 ERROR("Failed saving pending packet");
1263 goto dead;
1264 }
1265
1266 cd = (fr_channel_data_t *) lm;
1267 }
1268
1269 if (!s->blocked) {
1271 PERROR("Failed adding write callback to event loop");
1272 goto dead;
1273 }
1274
1275 s->blocked = true;
1276 }
1277
1278 s->pending = cd;
1279 return;
1280 }
1281
1282 PERROR("Failed writing to socket %s", s->listen->name);
1283
1284 switch (errno) {
1285 /*
1286 * As a special hack, check for something
1287 * that will never be returned from a
1288 * real write() routine. Which then
1289 * signals to us that we have to close
1290 * the socket, but NOT complain about it.
1291 */
1292 case ECONNREFUSED:
1293 case ECONNRESET:
1294 break;
1295
1296 /*
1297 * These are temporary errors. We just discard the data.
1298 */
1299 case ENETDOWN:
1300 case ENETUNREACH:
1301 if (li->app_io->error) li->app_io->error(li);
1302
1303 fr_message_done(&cd->m);
1304 return;
1305
1306 default:
1307 if (li->app_io->error) li->app_io->error(li);
1308 break;
1309 }
1310
1311 dead:
1312 fr_message_done(&cd->m);
1314 return;
1315 }
1316
1317 /*
1318 * If we've done a partial write, localize the message and continue.
1319 */
1320 if ((size_t) rcode < cd->m.data_size) {
1321 s->written = rcode;
1322 goto save_pending;
1323 }
1324
1325 discard:
1326 s->written = 0;
1327
1328 /*
1329 * Reset for the next message.
1330 */
1331 fr_message_done(&cd->m);
1332 nr->stats.out++;
1333 s->stats.out++;
1334
1335 /*
1336 * Grab the net entry.
1337 */
1338 fr_heap_pop((void **)&cd, &s->waiting);
1339 }
1340
1341 /*
1342 * We've successfully written all of the packets. Remove
1343 * the write callback.
1344 */
1346 PERROR("Failed removing write callback from event loop");
1348 }
1349
1350 s->blocked = false;
1351}
1352
1353/** Stop a socket's I/O events and close its descriptor
1354 *
1355 * Separate from freeing the socket because freeing it also frees s->ms, and
1356 * the workers hold messages allocated from that until they ack the channel
1357 * close. The descriptor can go as soon as we are done with it; the memory
1358 * cannot.
1359 *
1360 * Idempotent, so it is safe to call on a socket that is already dead.
1361 *
1362 * @param[in] s to close.
1363 * @return
1364 * - 0 on success.
1365 * - -1 if the I/O event could not be removed. The descriptor is closed
1366 * either way, so this is worth reporting but not worth stopping for.
1367 */
1369{
1370 int ret = 0;
1371
1372 if (s->closed) return 0;
1373 s->closed = true;
1374
1375 /*
1376 * NULL if the socket never made it into the event loop, which the
1377 * setup error paths rely on. This is the only place the event is
1378 * removed, so the handle cannot have gone stale under us.
1379 */
1380 if (s->ef) {
1381 ret = fr_event_fd_delete_handle(s->ef);
1382 s->ef = NULL;
1383 }
1384
1385 if (s->listen->app_io->close) {
1386 s->listen->app_io->close(s->listen);
1387 } else {
1388 close(s->listen->fd);
1389 }
1390
1391 return ret;
1392}
1393
1395{
1396 fr_network_t *nr = s->nr;
1398
1399 /*
1400 * Closing is a separate step, so that the descriptor can be
1401 * released while the message set stays alive for the workers.
1402 */
1403 fr_assert_msg(s->closed, "socket %d freed without being closed", s->number);
1404 fr_assert(s->outstanding == 0 || nr->exiting);
1405
1406 fr_rb_delete(nr->sockets, s);
1408
1409 if (s->pending) {
1411 s->pending = NULL;
1412 }
1413
1414 /*
1415 * Clean up any queued entries.
1416 */
1417 while ((fr_heap_pop((void **)&cd, &s->waiting) == 0) && cd) {
1418 fr_message_done(&cd->m);
1419 }
1420
1421 /* s->waiting is already talloc parented from s */
1422 talloc_free(s->listen);
1423
1424 return 0;
1425}
1426
1427
1428/** Handle a network control message callback for a new listener
1429 *
1430 * @param[in] data the message
1431 * @param[in] data_size size of the data
1432 * @param[in] now the current time
1433 * @param[in] uctx the network
1434 */
1435static void fr_network_listen_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
1436{
1437 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1438 fr_listen_t *li;
1439
1440 fr_assert(data_size == sizeof(li));
1441
1442 if (data_size != sizeof(li)) return;
1443
1444 li = talloc_get_type_abort(*((void * const *)data), fr_listen_t);
1445
1446 (void) fr_network_listen_add_self(nr, li);
1447}
1448
1449static void fr_network_limit_ringbuffer(fr_network_socket_t *s, int *num_messages_p, size_t *size_p)
1450{
1451 int num_messages;
1452 size_t size;
1453
1454 num_messages = s->listen->num_messages;
1455 if (num_messages < 8) num_messages = 8;
1456 if (num_messages > (1 << 20)) num_messages = (1 << 20);
1457
1458 size = s->listen->default_message_size * num_messages;
1459 if (size < (1 << 17)) size = (1 << 17);
1460 if (size > (100 * 1024 * 1024)) size = (100 * 1024 * 1024);
1461
1462 *num_messages_p = num_messages;
1463 *size_p = size;
1464}
1465
1467{
1469 fr_app_io_t const *app_io;
1470 size_t size;
1471 int num_messages;
1472
1473 fr_assert(li->app_io != NULL);
1474
1475 /*
1476 * Non-socket listeners just get told about the event
1477 * list, and nothing else.
1478 */
1479 if (li->non_socket_listener) {
1480 fr_assert(li->app_io->event_list_set != NULL);
1481 fr_assert(!li->app_io->read);
1482 fr_assert(!li->app_io->write);
1483
1484 li->app_io->event_list_set(li, nr->el, nr);
1485
1486 /*
1487 * We use fr_log() here to avoid the "Network - " prefix.
1488 */
1489 fr_log(nr->log, L_DBG, __FILE__, __LINE__, "Listener %s bound to virtual server %s",
1490 li->name, cf_section_name2(li->server_cs));
1491
1492 return 0;
1493 }
1494
1495 s = talloc_zero(nr, fr_network_socket_t);
1496 fr_assert(s != NULL);
1497 talloc_steal(s, li);
1498
1499 s->nr = nr;
1500 s->listen = li;
1501 s->number = nr->num_sockets++;
1502
1503 MEM(s->waiting = fr_heap_alloc(s, waiting_cmp, fr_channel_data_t, channel.heap_id, 0));
1504
1505 talloc_set_destructor(s, _network_socket_free);
1506
1507 /*
1508 * Put reasonable limits on the ring buffer size. Then
1509 * round it up to the nearest power of 2, which is
1510 * required by the ring buffer code.
1511 */
1512 fr_network_limit_ringbuffer(s, &num_messages, &size);
1513
1514 /*
1515 * Allocate the ring buffer for messages and packets.
1516 */
1517 s->ms = fr_message_set_create(s, num_messages,
1518 sizeof(fr_channel_data_t),
1519 size, false);
1520 if (!s->ms) {
1521 PERROR("Failed creating message buffers for network IO");
1522 if (network_socket_close(s) < 0) PWARN("Failed removing event for socket %d", s->number);
1523 talloc_free(s);
1524 return -1;
1525 }
1526
1527 app_io = s->listen->app_io;
1529
1530 if (fr_event_fd_insert(s, &s->ef, nr->el, s->listen->fd,
1534 s) < 0) {
1535 PERROR("Failed adding new socket to network event loop");
1536 if (network_socket_close(s) < 0) PWARN("Failed removing event for socket %d", s->number);
1537 talloc_free(s);
1538 return -1;
1539 }
1540
1541 /*
1542 * Start of with write updates being paused. We don't
1543 * care about being able to write if there's nothing to
1544 * write.
1545 */
1547
1548 /*
1549 * Add the listener before calling the app_io, so that
1550 * the app_io can find the listener which we're adding
1551 * here.
1552 */
1553 (void) fr_rb_insert(nr->sockets, s);
1554 (void) fr_rb_insert(nr->sockets_by_num, s);
1555
1556 if (app_io->event_list_set) app_io->event_list_set(s->listen, nr->el, nr);
1557
1558 /*
1559 * We use fr_log() here to avoid the "Network - " prefix.
1560 */
1561 fr_log(nr->log, L_DBG, __FILE__, __LINE__, "Listening on %s bound to virtual server %s",
1563
1564 DEBUG3("Using new socket %s with FD %d", s->listen->name, s->listen->fd);
1565
1566 return 0;
1567}
1568
1569/** Handle a network control message callback for a new "watch directory"
1570 *
1571 * @param[in] data the message
1572 * @param[in] data_size size of the data
1573 * @param[in] now the current time
1574 * @param[in] uctx the network
1575 */
1576static void fr_network_directory_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
1577{
1578 int num_messages;
1579 size_t size;
1580 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1581 fr_listen_t *li = talloc_get_type_abort(*((void * const *)data), fr_listen_t);
1583 fr_app_io_t const *app_io;
1585
1586 if (!fr_cond_assert(data_size == sizeof(li))) return;
1587
1588 memcpy(&li, data, sizeof(li));
1589
1590 s = talloc_zero(nr, fr_network_socket_t);
1591 fr_assert(s != NULL);
1592 talloc_steal(s, li);
1593
1594 s->nr = nr;
1595 s->listen = li;
1596 s->number = nr->num_sockets++;
1597
1598 MEM(s->waiting = fr_heap_alloc(s, waiting_cmp, fr_channel_data_t, channel.heap_id, 0));
1599
1600 talloc_set_destructor(s, _network_socket_free);
1601
1602 /*
1603 * Allocate the ring buffer for messages and packets.
1604 */
1605 fr_network_limit_ringbuffer(s, &num_messages, &size);
1606
1607 s->ms = fr_message_set_create(s, num_messages,
1608 sizeof(fr_channel_data_t),
1609 size, false);
1610 if (!s->ms) {
1611 PERROR("Failed creating message buffers for directory IO");
1612 if (network_socket_close(s) < 0) PWARN("Failed removing event for socket %d", s->number);
1613 talloc_free(s);
1614 return;
1615 }
1616
1617 app_io = s->listen->app_io;
1618
1619 if (app_io->event_list_set) app_io->event_list_set(s->listen, nr->el, nr);
1620
1622
1623 if (fr_event_filter_insert(s, &s->ef, nr->el, s->listen->fd, s->filter,
1624 &funcs,
1625 app_io->error ? fr_network_error : NULL,
1626 s) < 0) {
1627 PERROR("Failed adding directory monitor event loop");
1628 if (network_socket_close(s) < 0) PWARN("Failed removing event for socket %d", s->number);
1629 talloc_free(s);
1630 return;
1631 }
1632
1633 (void) fr_rb_insert(nr->sockets, s);
1634 (void) fr_rb_insert(nr->sockets_by_num, s);
1635
1636 DEBUG3("Using new socket with FD %d", s->listen->fd);
1637}
1638
1639/** Handle a network control message callback for a new worker
1640 *
1641 * @param[in] data the message
1642 * @param[in] data_size size of the data
1643 * @param[in] now the current time
1644 * @param[in] uctx the network
1645 */
1646static void fr_network_worker_started_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
1647{
1648 int i;
1649 fr_network_t *nr = uctx;
1650 fr_worker_t *worker;
1652 fr_worker_config_t const *worker_config;
1653
1654 fr_assert(data_size == sizeof(worker));
1655
1656 if (nr->num_workers >= nr->max_workers) {
1657 ERROR("Too many workers");
1658 return;
1659 }
1660
1661 memcpy(&worker, data, data_size);
1662 (void) talloc_get_type_abort(worker, fr_worker_t);
1663 worker_config = fr_worker_config(worker);
1664
1665 MEM(w = talloc_zero(nr, fr_network_worker_t));
1666
1667 w->worker = worker;
1668 w->reply_ms = fr_message_set_create(w, worker_config->message_set_size, sizeof(fr_channel_data_t),
1669 worker_config->ring_buffer_size, false);
1670 fr_fatal_assert_msg(w->reply_ms, "Failed creating reply message set");
1671 w->channel = fr_worker_channel_create(worker, w, nr->control, w->reply_ms);
1673 fr_fatal_assert_msg(w->channel, "Failed creating new channel");
1674
1677
1678 /*
1679 * Insert the worker into the array of workers.
1680 */
1681 for (i = 0; i < nr->max_workers; i++) {
1682 if (nr->workers[i]) continue;
1683
1684 nr->workers[i] = w;
1685 nr->num_workers++;
1686 return;
1687 }
1688}
1689
1690/** Handle a network control message callback for a packet sent to a socket
1691 *
1692 * @param[in] data the message
1693 * @param[in] data_size size of the data
1694 * @param[in] now the current time
1695 * @param[in] uctx the network
1696 */
1697static void fr_network_inject_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
1698{
1699 fr_network_t *nr = uctx;
1700 fr_network_inject_t my_inject;
1702
1703 fr_assert(data_size == sizeof(my_inject));
1704
1705 memcpy(&my_inject, data, data_size);
1706 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = my_inject.listen });
1707 if (!s) {
1708 talloc_free(my_inject.packet); /* MUST be it's own TALLOC_CTX */
1709 return;
1710 }
1711
1712 /*
1713 * Inject the packet, and then read it back from the
1714 * network.
1715 */
1716 if (s->listen->app_io->inject(s->listen, my_inject.packet, my_inject.packet_len, my_inject.recv_time) == 0) {
1717 fr_network_read(nr->el, s->listen->fd, 0, s);
1718 }
1719
1720 talloc_free(my_inject.packet);
1721}
1722
1723/** Run the event loop 'pre' callback
1724 *
1725 * This function MUST DO NO WORK. All it does is check if there's
1726 * work, and tell the event code to return to the main loop if
1727 * there's work to do.
1728 *
1729 * @param[in] now the current time.
1730 * @param[in] wake the time when the event loop will wake up.
1731 * @param[in] uctx the network
1732 */
1734{
1735 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1736
1737 if (fr_heap_num_elements(nr->replies) > 0) return 1;
1738
1739 return 0;
1740}
1741
1742/** Handle replies after all FD and timer events have been serviced
1743 *
1744 * @param el the event loop
1745 * @param now the current time (mostly)
1746 * @param uctx the fr_network_t
1747 */
1749{
1751 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1752
1753 /*
1754 * Pull the replies off of our global heap, and try to
1755 * push them to the individual sockets.
1756 */
1757 while ((fr_heap_pop((void **)&cd, &nr->replies) == 0) && cd) {
1758 fr_listen_t *li;
1760
1761 li = cd->listen;
1762
1763 /*
1764 * @todo - cache this somewhere so we don't need
1765 * to do an rbtree lookup for every packet.
1766 */
1767 fr_rb_find((void **)&s, nr->sockets, &(fr_network_socket_t){ .listen = li });
1768
1769 /*
1770 * This shouldn't happen, but be safe...
1771 */
1772 if (!s) {
1773 fr_message_done(&cd->m);
1774 continue;
1775 }
1776
1777 if (cd->m.status != FR_MESSAGE_LOCALIZED) {
1778 fr_assert(s->outstanding > 0);
1779 s->outstanding--;
1780 }
1781
1782 /*
1783 * Just mark the message done, and skip it.
1784 */
1785 if (s->dead) {
1786 fr_message_done(&cd->m);
1787
1788 /*
1789 * No more packets, it's safe to delete
1790 * the socket.
1791 */
1792 if (!s->outstanding) talloc_free(s);
1793
1794 continue;
1795 }
1796
1797 /*
1798 * No data to write to the socket, so we skip the message.
1799 */
1800 if (!cd->m.data_size) {
1801 fr_message_done(&cd->m);
1802 continue;
1803 }
1804
1805 (void) fr_heap_insert(&s->waiting, cd);
1806
1807 /*
1808 * No pending message, write it. If there is a pending write, the message will be left
1809 * in the waiting queue.
1810 */
1811 if (!s->pending) {
1812 fr_assert(!s->blocked);
1813 fr_network_write(nr->el, s->listen->fd, 0, s);
1814 }
1815 }
1816}
1817
1819{
1820 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1821
1822 return (nr->num_workers == 0) ? 1 : 0;
1823}
1824
1826{
1827 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1828
1829 if (nr->num_workers == 0) fr_event_loop_exit(el, SIGINT);
1830}
1831
1832/** Add events to the loop which will exit when the number of workers has reached zero
1833 */
1835{
1837 fr_strerror_const("Failed adding close pre-event to event list");
1838 return -1;
1839 }
1841 fr_strerror_const("Failed adding close post-event to event list");
1842 return -1;
1843 }
1844 return 0;
1845}
1846
1847/** Remove "close" events from the event loop.
1848 */
1854
1855/** Stop a network thread in an orderly way
1856 *
1857 * @param[in] nr the network to stop
1858 */
1860{
1862 int ret = 0;
1863
1864 (void) talloc_get_type_abort(nr, fr_network_t);
1865
1866 /*
1867 * Close the network sockets, but leave them allocated.
1868 *
1869 * Freeing a socket frees its message set, and the workers hold
1870 * messages allocated from it until they ack the channel close
1871 * signalled below. fr_network() frees them after its loop, which
1872 * already runs until every worker has acked.
1873 */
1874 {
1875 fr_network_socket_t **sockets;
1876 size_t len;
1877 size_t i;
1878
1879 if (fr_rb_flatten_inorder(nr, (void ***)&sockets, nr->sockets) < 0) return -1;
1880 len = talloc_array_length(sockets);
1881
1882 for (i = 0; i < len; i++) {
1883 if (network_socket_close(sockets[i]) < 0) {
1884 PWARN("Failed removing event for socket %d", sockets[i]->number);
1885 ret = -1;
1886 }
1887 }
1888
1889 talloc_free(sockets);
1890 }
1891
1892
1893 /*
1894 * Clean up all outstanding replies.
1895 *
1896 * We can't do this after signalling the
1897 * workers to close, because they free
1898 * their message sets, and we end up
1899 * getting random use-after-free errors
1900 * as there's a race between the network
1901 * popping replies, and the workers
1902 * freeing their message sets.
1903 *
1904 * This isn't perfect, and we might still
1905 * lose some replies, but it's good enough
1906 * for now.
1907 *
1908 * @todo - call transport "done" for the reply, so that
1909 * it knows the replies are done, too.
1910 */
1911 while ((fr_heap_pop((void **)&cd, &nr->replies) == 0) && cd) {
1912 fr_message_done(&cd->m);
1913 }
1914
1915 /*
1916 * Nothing may be queued from here on: a worker that has been
1917 * signalled discards whatever arrives late, without a reply.
1918 */
1919 nr->exiting = true;
1920
1921 /*
1922 * Signal the workers that we're closing
1923 *
1924 * nr->num_workers is decremented every
1925 * time a worker closes a socket.
1926 *
1927 * When nr->num_workers == 0, the event
1928 * loop (fr_network()) will exit.
1929 */
1930 {
1931 int i;
1932
1933 for (i = 0; i < nr->num_workers; i++) {
1934 fr_network_worker_t *worker = nr->workers[i];
1935
1937 }
1938 }
1939
1943
1944 return ret;
1945}
1946
1947/** Read handler for signal pipe
1948 *
1949 */
1950static void _signal_pipe_read(UNUSED fr_event_list_t *el, int fd, UNUSED int flags, void *uctx)
1951{
1952 fr_network_t *nr = talloc_get_type_abort(uctx, fr_network_t);
1953 uint8_t buff;
1954
1955 if (read(fd, &buff, sizeof(buff)) < 0) {
1956 ERROR("Failed reading signal - %s", fr_syserror(errno));
1957 return;
1958 }
1959
1960 fr_assert(buff == 1);
1961
1962 /*
1963 * fr_network_stop() will signal the workers
1964 * to exit (by closing their channels).
1965 *
1966 * When we get the ack, we decrement our
1967 * nr->num_workers counter.
1968 *
1969 * When the counter reaches 0, the event loop
1970 * exits.
1971 */
1972 DEBUG2("Signalled to exit");
1973
1974 if (unlikely(fr_network_destroy(nr) < 0)) {
1975 PERROR("Failed stopping network");
1976 }
1977}
1978
1979/** The main network worker function.
1980 *
1981 * @param[in] nr the network data structure to run.
1982 */
1984{
1985 /*
1986 * Run until we're told to exit AND the number of
1987 * workers has dropped to zero.
1988 *
1989 * This is important as if we exit too early we
1990 * free the channels out from underneath the
1991 * workers and they read uninitialised memory.
1992 *
1993 * Whenever a worker ACKs our close notification
1994 * nr->num_workers is decremented, so when
1995 * nr->num_workers == 0, all workers have ACKd
1996 * our close and are no longer using the channel.
1997 */
1998 while (likely(!(nr->exiting && (nr->num_workers == 0)))) {
1999 bool wait_for_event;
2000 int num_events;
2001
2002 /*
2003 * There are runnable requests. We still service
2004 * the event loop, but we don't wait for events.
2005 */
2006 wait_for_event = (fr_heap_num_elements(nr->replies) == 0);
2007
2008 /*
2009 * Check the event list. If there's an error
2010 * (e.g. exit), we stop looping and clean up.
2011 */
2012 DEBUG4("Gathering events - %s", wait_for_event ? "will wait" : "Will not wait");
2013 num_events = fr_event_corral(nr->el, fr_time(), wait_for_event);
2014 DEBUG4("%u event(s) pending%s",
2015 num_events == -1 ? 0 : num_events, num_events == -1 ? " - event loop exiting" : "");
2016 if (num_events < 0) break;
2017
2018 /*
2019 * Service outstanding events.
2020 */
2021 if (num_events > 0) {
2022 DEBUG4("Servicing event(s)");
2023 fr_event_service(nr->el);
2024 }
2025 }
2026
2027 /*
2028 * Free the sockets. The loop above only exits once every worker
2029 * has acked, so the messages allocated from the sockets' message
2030 * sets have all been returned, and fr_network_destroy() left the
2031 * sockets closed but allocated for exactly this point. It cannot do
2032 * this itself, as it returns long before the first ack arrives. The
2033 * error exit above lands here too, where no ack is ever coming.
2034 */
2035 {
2036 fr_network_socket_t **sockets;
2037 size_t len;
2038 size_t i;
2039
2040 if (fr_rb_flatten_inorder(nr, (void ***)&sockets, nr->sockets) < 0) {
2041 PWARN("Failed enumerating sockets, leaving them to be freed with the network");
2042 return;
2043 }
2044 len = talloc_array_length(sockets);
2045
2046 for (i = 0; i < len; i++) talloc_free(sockets[i]);
2047
2048 talloc_free(sockets);
2049 }
2050}
2051
2052/** Signal a network thread to exit
2053 *
2054 * @note Request to exit will be processed asynchronously.
2055 *
2056 * @param[in] nr the network data structure to manage
2057 * @return
2058 * - 0 on success.
2059 * - -1 on failure.
2060 */
2062{
2063 if (write(nr->signal_pipe[1], &(uint8_t){ 0x01 }, 1) < 0) {
2064 fr_strerror_printf("Failed signalling network thread to exit - %s", fr_syserror(errno));
2065 return -1;
2066 }
2067
2068 return 0;
2069}
2070
2071/** Free any resources associated with a network thread
2072 *
2073 */
2075{
2076 if (nr->signal_pipe[0] >= 0) close(nr->signal_pipe[0]);
2077 if (nr->signal_pipe[1] >= 0) close(nr->signal_pipe[1]);
2078
2079 return 0;
2080}
2081
2082/** Create a network
2083 *
2084 * @param[in] ctx The talloc ctx
2085 * @param[in] el The event list
2086 * @param[in] name Networker identifier.
2087 * @param[in] logger The destination for all logging messages
2088 * @param[in] lvl Log level
2089 * @param[in] config configuration structure.
2090 * @return
2091 * - NULL on error
2092 * - fr_network_t on success
2093 */
2094fr_network_t *fr_network_create(TALLOC_CTX *ctx, fr_event_list_t *el, char const *name,
2095 fr_log_t const *logger, fr_log_lvl_t lvl,
2097{
2098 fr_network_t *nr;
2099
2100 nr = talloc_zero(ctx, fr_network_t);
2101 if (!nr) {
2102 fr_strerror_const("Failed allocating memory");
2103 return NULL;
2104 }
2105 talloc_set_destructor(nr, _fr_network_free);
2106
2107 nr->name = talloc_strdup(nr, name);
2108
2109 nr->thread_id = pthread_self();
2110 nr->el = el;
2111 nr->log = logger;
2112 nr->lvl = lvl;
2113
2115 nr->num_workers = 0;
2116 nr->signal_pipe[0] = -1;
2117 nr->signal_pipe[1] = -1;
2118 if (config) nr->config = *config;
2119
2120 nr->aq_control = fr_atomic_queue_talloc(nr, 1024);
2121 if (!nr->aq_control) {
2122 talloc_free(nr);
2123 return NULL;
2124 }
2125
2126 nr->control = fr_control_create(nr, el, nr->aq_control, 6);
2127 if (!nr->control) {
2128 fr_strerror_const_push("Failed creating control queue");
2129 fail:
2130 talloc_free(nr);
2131 return NULL;
2132 }
2133
2134 /*
2135 * @todo - rely on thread-local variables. And then the
2136 * various users of this can check if (rb == nr->rb), and
2137 * if so, skip the whole control plane / kevent /
2138 * whatever roundabout thing.
2139 */
2141 if (!nr->rb) {
2142 fr_strerror_const_push("Failed creating ring buffer");
2143 fail2:
2144 talloc_free(nr->control);
2145 goto fail;
2146 }
2147
2149 fr_strerror_const_push("Failed adding channel callback");
2150 goto fail2;
2151 }
2152
2154 fr_strerror_const_push("Failed adding socket callback");
2155 goto fail2;
2156 }
2157
2159 fr_strerror_const_push("Failed adding socket callback");
2160 goto fail2;
2161 }
2162
2164 fr_strerror_const_push("Failed adding worker callback");
2165 goto fail2;
2166 }
2167
2169 fr_strerror_const_push("Failed adding packet injection callback");
2170 goto fail2;
2171 }
2172
2173 if (fr_control_open(nr->control) < 0) {
2174 fr_strerror_const_push("Failed opening control queue");
2175 goto fail2;
2176 }
2177
2178 /*
2179 * Create the various heaps.
2180 */
2182 if (!nr->sockets) {
2183 fr_strerror_const_push("Failed creating listen tree for sockets");
2184 goto fail2;
2185 }
2186
2188 if (!nr->sockets_by_num) {
2189 fr_strerror_const_push("Failed creating number tree for sockets");
2190 goto fail2;
2191 }
2192
2193 nr->replies = fr_heap_alloc(nr, reply_cmp, fr_channel_data_t, channel.heap_id, 0);
2194 if (!nr->replies) {
2195 fr_strerror_const_push("Failed creating heap for replies");
2196 goto fail2;
2197 }
2198
2199 if (fr_event_pre_insert(nr->el, fr_network_pre_event, nr) < 0) {
2200 fr_strerror_const("Failed adding pre-check to event list");
2201 goto fail2;
2202 }
2203
2204 if (fr_event_post_insert(nr->el, fr_network_post_event, nr) < 0) {
2205 fr_strerror_const("Failed inserting post-processing event");
2206 goto fail2;
2207 }
2208
2209 if (pipe(nr->signal_pipe) < 0) {
2210 fr_strerror_printf("Failed initialising signal pipe - %s", fr_syserror(errno));
2211 goto fail2;
2212 }
2213 if (fr_nonblock(nr->signal_pipe[0]) < 0) goto fail2;
2214 if (fr_nonblock(nr->signal_pipe[1]) < 0) goto fail2;
2215
2216 if (fr_event_fd_insert(nr, NULL, nr->el, nr->signal_pipe[0], _signal_pipe_read, NULL, NULL, nr) < 0) {
2217 fr_strerror_const("Failed inserting event for signal pipe");
2218 goto fail2;
2219 }
2220
2221 return nr;
2222}
2223
2224int fr_network_stats(fr_network_t const *nr, int num, uint64_t *stats)
2225{
2226 if (num < 0) return -1;
2227 if (num == 0) return 0;
2228
2229 stats[0] = nr->stats.in;
2230 if (num >= 2) stats[1] = nr->stats.out;
2231 if (num >= 3) stats[2] = nr->stats.dup;
2232 if (num >= 4) stats[3] = nr->stats.dropped;
2233 if (num >= 5) stats[4] = nr->num_workers;
2234
2235 if (num <= 5) return num;
2236
2237 return 5;
2238}
2239
2240void fr_network_stats_log(fr_network_t const *nr, fr_log_t const *log)
2241{
2242 int i;
2243
2244 /*
2245 * Dump all of the channel statistics.
2246 */
2247 for (i = 0; i < nr->max_workers; i++) {
2248 if (!nr->workers[i]) continue;
2249
2250 fr_channel_stats_log(nr->workers[i]->channel, log, __FILE__, __LINE__);
2251 }
2252}
2253
2254static int cmd_stats_self(FILE *fp, UNUSED FILE *fp_err, void *ctx, UNUSED fr_cmd_info_t const *info)
2255{
2256 fr_network_t const *nr = ctx;
2257
2258 fprintf(fp, "count.in\t%" PRIu64 "\n", nr->stats.in);
2259 fprintf(fp, "count.out\t%" PRIu64 "\n", nr->stats.out);
2260 fprintf(fp, "count.dup\t%" PRIu64 "\n", nr->stats.dup);
2261 fprintf(fp, "count.dropped\t%" PRIu64 "\n", nr->stats.dropped);
2262 fprintf(fp, "count.sockets\t%u\n", fr_rb_num_elements(nr->sockets));
2263
2264 return 0;
2265}
2266
2267static int cmd_socket_list(FILE *fp, UNUSED FILE *fp_err, void *ctx, UNUSED fr_cmd_info_t const *info)
2268{
2269 fr_network_t const *nr = ctx;
2272
2273 // @todo - note that this isn't thread-safe!
2274
2275 for (s = fr_rb_iter_init_inorder(nr->sockets, &iter);
2276 s != NULL;
2277 s = fr_rb_iter_next_inorder(nr->sockets, &iter)) {
2278 if (!s->listen->app_io->get_name) {
2279 fprintf(fp, "%s\n", s->listen->app_io->common.name);
2280 } else {
2281 fprintf(fp, "%d\t%s\n", s->number, s->listen->app_io->get_name(s->listen));
2282 }
2283 }
2284 return 0;
2285}
2286
2287static int cmd_stats_socket(FILE *fp, FILE *fp_err, void *ctx, fr_cmd_info_t const *info)
2288{
2289 fr_network_t const *nr = ctx;
2291
2292 fr_rb_find((void **)&s, nr->sockets_by_num, &(fr_network_socket_t){ .number = info->box[0]->vb_uint32 });
2293 if (!s) {
2294 fprintf(fp_err, "No such socket number '%s'.\n", info->argv[0]);
2295 return -1;
2296 }
2297
2298 fprintf(fp, "count.in\t%" PRIu64 "\n", s->stats.in);
2299 fprintf(fp, "count.out\t%" PRIu64 "\n", s->stats.out);
2300 fprintf(fp, "count.dup\t%" PRIu64 "\n", s->stats.dup);
2301 fprintf(fp, "count.dropped\t%" PRIu64 "\n", s->stats.dropped);
2302
2303 return 0;
2304}
2305
2306
2308 {
2309 .parent = "stats",
2310 .name = "network",
2311 .help = "Statistics for network threads.",
2312 .read_only = true
2313 },
2314
2315 {
2316 .parent = "stats network",
2317 .add_name = true,
2318 .name = "self",
2319 .func = cmd_stats_self,
2320 .help = "Show statistics for a specific network thread.",
2321 .read_only = true
2322 },
2323
2324 {
2325 .parent = "stats network",
2326 .add_name = true,
2327 .name = "socket",
2328 .syntax = "INTEGER",
2329 .func = cmd_stats_socket,
2330 .help = "Show statistics for a specific socket",
2331 .read_only = true
2332 },
2333
2334 {
2335 .parent = "show",
2336 .name = "network",
2337 .help = "Show information about network threads.",
2338 .read_only = true
2339 },
2340
2341 {
2342 .parent = "show network",
2343 .add_name = true,
2344 .name = "socket",
2345 .syntax = "list",
2346 .func = cmd_socket_list,
2347 .help = "List the sockets associated with this network thread.",
2348 .read_only = true
2349 },
2350
2352};
static int const char char buffer[256]
Definition acutest.h:576
fr_io_close_t close
Close the transport.
Definition app_io.h:60
fr_io_data_read_t read
Read from a socket to a data buffer.
Definition app_io.h:47
module_t common
Common fields to all loadable modules.
Definition app_io.h:34
fr_io_signal_t error
There was an error on the socket.
Definition app_io.h:59
fr_app_event_list_set_t event_list_set
Called by the network thread to pass an event list for use by the app_io_t.
Definition app_io.h:36
fr_io_data_inject_t inject
Inject a packet into a socket.
Definition app_io.h:50
fr_io_data_vnode_t vnode
Handle notifications that the VNODE has changed.
Definition app_io.h:52
fr_io_data_write_t write
Write from a data buffer to a socket.
Definition app_io.h:48
fr_io_name_t get_name
get the socket name
Definition app_io.h:70
Public structure describing an I/O path for a protocol.
Definition app_io.h:33
fr_app_priority_get_t priority
Assign a priority to the packet.
Definition application.h:90
#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 RCSID(id)
Definition build.h:560
#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
char const * cf_section_name2(CONF_SECTION const *cs)
Return the second identifier of a CONF_SECTION.
Definition cf_util.c:1363
void * fr_channel_requestor_uctx_get(fr_channel_t *ch)
Get network-specific data from a channel.
Definition channel.c:964
fr_table_num_sorted_t const channel_signals[]
Definition channel.c:152
bool fr_channel_recv_reply(fr_channel_t *ch)
Receive a reply message from the channel.
Definition channel.c:407
int fr_channel_signal_responder_close(fr_channel_t *ch)
Signal a responder that the channel is closing.
Definition channel.c:844
int fr_channel_send_request(fr_channel_t *ch, fr_channel_data_t *cd)
Send a request message into the channel.
Definition channel.c:305
int fr_channel_set_recv_reply(fr_channel_t *ch, fr_channel_recv_callback_t recv_reply, void *uctx)
Definition channel.c:972
void fr_channel_requestor_uctx_add(fr_channel_t *ch, void *uctx)
Add network-specific data to a channel.
Definition channel.c:952
void fr_channel_stats_log(fr_channel_t const *ch, fr_log_t const *log, char const *file, int line)
Definition channel.c:1006
fr_channel_event_t fr_channel_service_message(fr_time_t when, fr_channel_t **p_channel, void **uctx_out, void const *data, size_t data_size)
Service a control-plane message.
Definition channel.c:689
A full channel, which consists of two ends.
Definition channel.c:143
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
#define PRIORITY_NORMAL
Definition channel.h:153
#define PRIORITY_NOW
Definition channel.h:151
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
char const ** argv
text version of commands
Definition command.h:42
#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:177
#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:248
#define fr_cond_assert_msg(_x, _fmt,...)
Calls panic_action ifndef NDEBUG, else logs error and evaluates to value of _x.
Definition debug.h:194
#define fr_fatal_assert_msg(_x, _fmt,...)
Calls panic_action ifndef NDEBUG, else logs error and causes the server to exit immediately with code...
Definition debug.h:222
#define MEM(x)
Definition debug.h:38
#define ERROR(fmt,...)
Definition dhcpclient.c:40
static int sockfd
Definition dhcpclient.c:55
fr_dict_t const * fr_dict_proto_dict(fr_dict_t const *dict)
Definition dict_util.c:5413
#define fr_event_fd_insert(...)
Definition event.h:247
fr_event_filter_t
The type of filter to install for an FD.
Definition event.h:82
@ FR_EVENT_FILTER_VNODE
Filter for vnode subfilters.
Definition event.h:84
@ FR_EVENT_FILTER_IO
Combined filter for read/write functions/.
Definition event.h:83
#define fr_event_filter_update(...)
Definition event.h:239
#define fr_event_filter_insert(...)
Definition event.h:234
#define FR_EVENT_RESUME(_s, _f)
Re-add the filter for a func from kevent.
Definition event.h:131
#define FR_EVENT_SUSPEND(_s, _f)
Temporarily remove the filter for a func from kevent.
Definition event.h:115
fr_event_fd_cb_t extend
Additional files were added to a directory.
Definition event.h:198
Callbacks for the FR_EVENT_FILTER_IO filter.
Definition event.h:188
Structure describing a modification to a filter's state.
Definition event.h:96
Callbacks for the FR_EVENT_FILTER_VNODE filter.
Definition event.h:195
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
unsigned int fr_heap_index_t
Definition heap.h:82
#define fr_heap_alloc(_ctx, _cmp, _type, _field, _init)
Creates a heap that can be used with non-talloced elements.
Definition heap.h:102
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_INDEX_INVALID
Definition heap.h:85
The main heap structure.
Definition heap.h:68
talloc_free(hp)
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
static void fr_network_inject_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
Handle a network control message callback for a packet sent to a socket.
Definition network.c:1697
fr_ring_buffer_t * rb
ring buffer for my control-plane messages
Definition network.c:125
static int fr_network_pre_close_event(UNUSED fr_time_t now, UNUSED fr_time_delta_t wake, void *uctx)
Definition network.c:1818
fr_cmd_table_t cmd_network_table[]
Definition network.c:2307
size_t written
however much we did in a partial write
Definition network.c:91
int fr_network_listen_send_packet(fr_network_t *nr, fr_listen_t *parent, fr_listen_t *li, const uint8_t *buffer, size_t buflen, fr_time_t recv_time, void *packet_ctx)
Send a packet to the worker.
Definition network.c:816
fr_atomic_queue_t * aq_control
atomic queue for control messages sent to me
Definition network.c:121
static int cmd_stats_socket(FILE *fp, FILE *fp_err, void *ctx, fr_cmd_info_t const *info)
Definition network.c:2287
int fr_network_listen_add(fr_network_t *nr, fr_listen_t *li)
Add a fr_listen_t to a network.
Definition network.c:236
bool suspended
whether or not we're suspended.
Definition network.c:116
void fr_network_close_event_delete(fr_network_t *nr)
Remove "close" events from the event loop.
Definition network.c:1849
int fr_network_worker_add(fr_network_t *nr, fr_worker_t *worker)
Add a worker to a network in a different thread.
Definition network.c:305
int fr_network_destroy(fr_network_t *nr)
Stop a network thread in an orderly way.
Definition network.c:1859
fr_network_t * nr
O(N) issues in talloc.
Definition network.c:73
fr_io_stats_t stats
Definition network.c:66
fr_listen_t * listen
Definition network.c:47
static fr_cmp_ret_t socket_num_cmp(void const *one, void const *two)
Definition network.c:186
fr_event_fd_t * ef
the I/O event, if it was inserted.
Definition network.c:78
uint8_t * packet
Definition network.c:48
static int cmd_stats_self(FILE *fp, UNUSED FILE *fp_err, void *ctx, UNUSED fr_cmd_info_t const *info)
Definition network.c:2254
fr_log_t const * log
log destination
Definition network.c:118
int fr_network_directory_add(fr_network_t *nr, fr_listen_t *li)
Add a "watch directory" call to a network.
Definition network.c:290
static int _fr_network_free(fr_network_t *nr)
Free any resources associated with a network thread.
Definition network.c:2074
fr_heap_index_t heap_id
for the sockets_by_num heap
Definition network.c:75
#define RTT(_old, _new)
Definition network.c:512
void fr_network(fr_network_t *nr)
The main network worker function.
Definition network.c:1983
fr_message_set_t * ms
message buffers for this socket.
Definition network.c:88
int fr_network_listen_delete(fr_network_t *nr, fr_listen_t *li)
Delete a socket from a network.
Definition network.c:271
int num_blocked
number of blocked workers
Definition network.c:137
ssize_t fr_network_listen_outstanding(fr_network_t *nr, fr_listen_t *li)
Get the number of outstanding packets.
Definition network.c:858
char const * name
Network ID for logging.
Definition network.c:112
void fr_network_worker_add_self(fr_network_t *nr, fr_worker_t *worker)
Add a worker to a network in the same thread.
Definition network.c:325
unsigned int outstanding
number of outstanding packets sent to the worker
Definition network.c:85
static _Thread_local fr_ring_buffer_t * fr_network_rb
Definition network.c:44
int number
unique ID
Definition network.c:74
static fr_event_update_t const resume_write[]
Definition network.c:476
fr_message_set_t * reply_ms
message set the worker will use to send reply data.
Definition network.c:64
fr_time_delta_t predicted
predicted processing time for one packet
Definition network.c:59
static int fr_network_pre_event(fr_time_t now, fr_time_delta_t wake, void *uctx)
static fr_event_update_t const pause_read[]
Definition network.c:461
fr_worker_t * worker
worker pointer
Definition network.c:65
int fr_network_sendto_worker(fr_network_t *nr, fr_listen_t *li, void *packet_ctx, uint8_t const *data, size_t data_len, fr_time_t recv_time)
Definition network.c:1110
int fr_network_exit(fr_network_t *nr)
Signal a network thread to exit.
Definition network.c:2061
#define MAX_WORKERS
Definition network.c:42
int fr_network_listen_inject(fr_network_t *nr, fr_listen_t *li, uint8_t const *packet, size_t packet_len, fr_time_t recv_time)
Inject a packet for a listener to read.
Definition network.c:410
fr_listen_t * listen
I/O ctx and functions.
Definition network.c:86
int num_sockets
actually a counter...
Definition network.c:140
fr_rb_node_t listen_node
rbtree node for looking up by listener.
Definition network.c:70
static void fr_network_channel_callback(void const *data, size_t data_size, fr_time_t now, void *uctx)
Handle a network control message callback for a channel.
Definition network.c:568
static void fr_network_vnode_extend(UNUSED fr_event_list_t *el, int sockfd, int fflags, void *ctx)
Get a notification that a vnode changed.
Definition network.c:1154
static void _signal_pipe_read(UNUSED fr_event_list_t *el, int fd, UNUSED int flags, void *uctx)
Read handler for signal pipe.
Definition network.c:1950
#define OUTSTANDING(_x)
Definition network.c:629
int num_workers
number of active workers
Definition network.c:136
int num_pending_workers
number of workers we're waiting to start.
Definition network.c:138
fr_rb_tree_t * sockets
list of sockets we're managing, ordered by the listener
Definition network.c:133
bool closed
the descriptor has been closed.
Definition network.c:83
pthread_t thread_id
for self
Definition network.c:114
fr_log_lvl_t lvl
debug log level
Definition network.c:119
int signal_pipe[2]
Pipe for signalling the worker in an orderly way.
Definition network.c:142
fr_channel_data_t * pending
the currently pending partial packet
Definition network.c:93
static void fr_network_write(UNUSED fr_event_list_t *el, UNUSED int sockfd, UNUSED int flags, void *ctx)
Write packets to the network.
Definition network.c:1206
fr_event_list_t * el
our event list
Definition network.c:127
fr_heap_t * replies
replies from the worker, ordered by priority / origin time
Definition network.c:129
static int fr_network_send_request(fr_network_t *nr, fr_channel_data_t *cd)
Send a message on the "best" channel.
Definition network.c:636
static fr_event_update_t const resume_read[]
Definition network.c:466
void fr_network_stats_log(fr_network_t const *nr, fr_log_t const *log)
Definition network.c:2240
fr_heap_t * waiting
packets waiting to be written
Definition network.c:94
int fr_network_stats(fr_network_t const *nr, int num, uint64_t *stats)
Definition network.c:2224
fr_heap_index_t heap_id
workers are in a heap
Definition network.c:57
static void fr_network_worker_started_callback(void const *data, size_t data_size, fr_time_t now, void *uctx)
bool blocked
is this worker blocked?
Definition network.c:61
static void fr_network_read(UNUSED fr_event_list_t *el, int sockfd, UNUSED int flags, void *ctx)
Read a packet from the network.
Definition network.c:919
static bool is_network_thread(fr_network_t const *nr)
Definition network.c:224
static fr_cmp_ret_t reply_cmp(void const *one, void const *two)
Definition network.c:157
fr_rb_node_t num_node
rbtree node for looking up by number.
Definition network.c:71
static void fr_network_error(UNUSED fr_event_list_t *el, UNUSED int sockfd, int flags, int fd_errno, void *ctx)
Handle errors for a socket.
Definition network.c:1179
fr_io_stats_t stats
Definition network.c:95
static fr_ring_buffer_t * fr_network_rb_init(void)
Initialise thread local storage.
Definition network.c:206
fr_channel_data_t * cd
cached in case of allocation & read error
Definition network.c:89
static int fr_network_listen_add_self(fr_network_t *nr, fr_listen_t *listen)
Definition network.c:1466
static void fr_network_suspend(fr_network_t *nr)
Definition network.c:481
bool dead
is it dead?
Definition network.c:81
size_t leftover
leftover data from a previous read
Definition network.c:90
int fr_network_close_event_insert(fr_network_t *nr)
Add events to the loop which will exit when the number of workers has reached zero.
Definition network.c:1834
static void fr_network_post_event(fr_event_list_t *el, fr_time_t now, void *uctx)
fr_network_worker_t * workers[MAX_WORKERS]
each worker
Definition network.c:148
fr_time_t recv_time
Definition network.c:50
static void fr_network_unsuspend(fr_network_t *nr)
Definition network.c:496
fr_rb_tree_t * sockets_by_num
ordered by number;
Definition network.c:134
fr_network_config_t config
configuration
Definition network.c:147
void fr_network_listen_read(fr_network_t *nr, fr_listen_t *li)
Signal the network to read from a listener.
Definition network.c:336
static fr_cmp_ret_t waiting_cmp(void const *one, void const *two)
Definition network.c:168
fr_io_stats_t stats
Definition network.c:131
static void fr_network_limit_ringbuffer(fr_network_socket_t *s, int *num_messages_p, size_t *size_p)
Definition network.c:1449
static int _fr_network_rb_free(void *arg)
Definition network.c:197
int max_workers
maximum number of allowed workers
Definition network.c:139
void fr_network_listen_write(fr_network_t *nr, fr_listen_t *li, uint8_t const *packet, size_t packet_len, void *packet_ctx, fr_time_t request_time)
Inject a packet for a listener to write.
Definition network.c:362
static void fr_network_recv_reply(fr_channel_t *ch, fr_channel_data_t *cd, void *uctx)
Callback which handles a message being received on the network side.
Definition network.c:520
bool exiting
are we exiting?
Definition network.c:145
fr_event_filter_t filter
what type of filter it is
Definition network.c:77
static void fr_network_socket_dead(fr_network_t *nr, fr_network_socket_t *s)
Definition network.c:874
fr_channel_t * channel
channel to the worker
Definition network.c:63
static fr_event_update_t const pause_write[]
Definition network.c:471
fr_network_t * fr_network_create(TALLOC_CTX *ctx, fr_event_list_t *el, char const *name, fr_log_t const *logger, fr_log_lvl_t lvl, fr_network_config_t const *config)
Create a network.
Definition network.c:2094
static int _network_socket_free(fr_network_socket_t *s)
Definition network.c:1394
static void fr_network_directory_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
Handle a network control message callback for a new "watch directory".
Definition network.c:1576
static int network_socket_close(fr_network_socket_t *s)
Stop a socket's I/O events and close its descriptor.
Definition network.c:1368
static fr_cmp_ret_t socket_listen_cmp(void const *one, void const *two)
Definition network.c:179
static int cmd_socket_list(FILE *fp, UNUSED FILE *fp_err, void *ctx, UNUSED fr_cmd_info_t const *info)
Definition network.c:2267
fr_time_delta_t cpu_time
how much CPU time this worker has spent
Definition network.c:58
fr_control_t * control
the control plane
Definition network.c:123
bool blocked
is it blocked?
Definition network.c:82
static void fr_network_listen_callback(void const *data, size_t data_size, UNUSED fr_time_t now, void *uctx)
Handle a network control message callback for a new listener.
Definition network.c:1435
static void fr_network_post_close_event(fr_event_list_t *el, UNUSED fr_time_t now, void *uctx)
Definition network.c:1825
Associate a worker thread with a network thread.
Definition network.c:56
uint32_t max_outstanding
Definition network.h:46
#define PERROR(_fmt,...)
Definition log.h:233
#define DEBUG3(_fmt,...)
Definition log.h:271
#define PWARN(_fmt,...)
Definition log.h:232
#define DEBUG4(_fmt,...)
Definition log.h:272
#define HEXDUMP2(_data, _len, _fmt,...)
Definition log.h:739
#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:2205
int fr_event_pre_delete(fr_event_list_t *el, fr_event_status_cb_t callback, void *uctx)
Delete a pre-event callback from the event list.
Definition event.c:2003
int fr_event_fd_delete_handle(fr_event_fd_t *ef)
Remove a file descriptor from the event loop, by handle.
Definition event.c:1228
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:2073
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:2051
#define fr_time()
Definition event.c:60
int fr_event_pre_insert(fr_event_list_t *el, fr_event_status_cb_t callback, void *uctx)
Add a pre-event callback to the event list.
Definition event.c:1981
void fr_event_loop_exit(fr_event_list_t *el, int code)
Signal an event loop exit with the specified code.
Definition event.c:2383
int fr_event_fd_delete(fr_event_list_t *el, int fd, fr_event_filter_t filter)
Remove a file descriptor from the event loop.
Definition event.c:1203
int fr_event_post_insert(fr_event_list_t *el, fr_event_post_cb_t callback, void *uctx)
Add a post-event callback to the event list.
Definition event.c:2028
A file descriptor/filter event.
Definition event.c:260
Stores all information relating to an event list.
Definition event.c:377
#define fr_log(_log, _lvl, _file, _line, _fmt,...)
Definition log.h:172
fr_log_lvl_t
Definition log.h:64
@ L_DBG
Only displayed when debugging is enabled.
Definition log.h:56
bool read_hexdump
Do we debug hexdump packets as they're read.
Definition listen.h:53
size_t num_messages
for the message ring buffer
Definition listen.h:57
bool non_socket_listener
special internal listener that does not use sockets.
Definition listen.h:47
char const * name
printable name for this socket - set by open
Definition listen.h:29
void const * app_instance
Definition listen.h:39
size_t default_message_size
copied from app_io, but may be changed
Definition listen.h:56
bool write_hexdump
Do we debug hexdump packets as they're written.
Definition listen.h:54
fr_app_t const * app
Definition listen.h:38
CONF_SECTION * server_cs
CONF_SECTION of the server.
Definition listen.h:42
bool no_write_callback
sometimes we don't need to do writes
Definition listen.h:46
int fd
file descriptor for this socket - set by open
Definition listen.h:28
fr_dict_t const * dict
dictionary for this listener
Definition listen.h:30
bool needs_full_setup
Set to true to avoid the short cut when adding the listener.
Definition listen.h:48
fr_app_io_t const * app_io
I/O path functions.
Definition listen.h:32
unsigned int uint32_t
long int ssize_t
unsigned char uint8_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
void fr_message_and_data_reset(fr_message_set_t *ms, fr_message_t *m)
Cancel a reservation made by fr_message_and_data_reserve(), returning the slot to the set.
Definition message.c:1026
int fr_message_done(fr_message_t *m)
Mark a message as done.
Definition message.c:195
fr_message_t * fr_message_localize(TALLOC_CTX *ctx, fr_message_t *m, size_t message_size)
Localize a message by copying it to local storage.
Definition message.c:247
fr_message_t * fr_message_and_data_commit_with_leftover(fr_message_set_t *ms, fr_message_t *m, size_t actual_packet_size, size_t leftover, size_t reserve_size)
Allocate packet data for a message, and reserve a new message.
Definition message.c:1159
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
fr_message_t * fr_message_and_data_alloc(fr_message_set_t *ms, size_t size)
Reserve and commit a message atomically.
Definition message.c:1108
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_MESSAGE_USED
Definition message.h:39
@ FR_MESSAGE_LOCALIZED
Definition message.h:40
fr_message_status_t status
free, used, done, etc.
Definition message.h:45
int fr_nonblock(UNUSED int fd)
Definition misc.c:293
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 DEBUG2(fmt,...)
static fr_app_io_t app_io
uint32_t fr_rand(void)
Return a 32-bit random number.
Definition rand.c:104
uint32_t fr_rb_num_elements(fr_rb_tree_t *tree)
Return how many nodes there are in a tree.
Definition rb.c:807
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
void * fr_rb_iter_init_inorder(fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Initialise an in-order iterator.
Definition rb.c:850
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
void * fr_rb_iter_next_inorder(UNUSED fr_rb_tree_t *tree, fr_rb_iter_inorder_t *iter)
Return the next node.
Definition rb.c:876
#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
int fr_rb_flatten_inorder(TALLOC_CTX *ctx, void **out[], fr_rb_tree_t *tree)
Iterator structure for in-order traversal of an rbtree.
Definition rb.h:319
The main red black tree structure.
Definition rb.h:71
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
static char buff[sizeof("18446744073709551615")+3]
Definition size_tests.c:37
Definition log.h:93
char const * fr_syserror(int num)
Guaranteed to be thread-safe version of strerror.
Definition syserror.c:243
#define fr_table_str_by_value(_table, _number, _def)
Convert an integer to a string.
Definition table.h:804
#define talloc_get_type_abort_const
Definition talloc.h:117
#define talloc_strdup(_ctx, _str)
Definition talloc.h:149
static fr_time_delta_t fr_time_delta_from_msec(int64_t msec)
Definition time.h:575
static fr_time_delta_t fr_time_delta_add(fr_time_delta_t a, fr_time_delta_t b)
Definition time.h:255
#define fr_time_delta_lt(_a, _b)
Definition time.h:285
#define fr_time_wrap(_time)
Definition time.h:145
#define fr_time_delta_ispos(_a)
Definition time.h:290
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
static fr_event_list_t * el
static fr_slen_t parent
Definition pair.h:860
void fr_perror(char const *fmt,...)
Print the current error to stderr with a prefix.
Definition strerror.c:737
#define fr_strerror_printf(_fmt,...)
Log to thread local error buffer.
Definition strerror.h:64
#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:1367
static fr_sbuff_err_t char size_t * len
Definition value.h:1062
fr_dict_t const * virtual_server_dict_by_cs(CONF_SECTION const *cs)
Return the namespace for specified CONF_SECTION.
fr_channel_t * fr_worker_channel_create(fr_worker_t *worker, TALLOC_CTX *ctx, fr_control_t *master, void *uctx)
Create a channel to the worker.
Definition worker.c:1728
int fr_worker_listen_cancel(fr_worker_t *worker, fr_listen_t const *li)
Definition worker.c:1754
fr_worker_config_t const * fr_worker_config(fr_worker_t *worker)
Definition worker.c:1820
A worker which takes packets from a master, and processes them.
Definition worker.c:91
int message_set_size
default start number of messages
Definition worker.h:73
#define FR_CONTROL_ID_INJECT
Definition worker.h:39
#define FR_CONTROL_ID_DIRECTORY
Definition worker.h:38
#define FR_CONTROL_ID_LISTEN
Definition worker.h:36
#define FR_CONTROL_ID_WORKER
Definition worker.h:37
int ring_buffer_size
default start size for the ring buffers
Definition worker.h:74