The FreeRADIUS server $Id: f3670dba8951ca10eb4948feb3dc3db9423a334f $
Loading...
Searching...
No Matches
channel.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: b45543b357714504d5ca0ceeab7a37f931a36fee $
19 *
20 * @brief Two-way thread-safe channels.
21 * @file io/channel.c
22 *
23 * @copyright 2016 Alan DeKok (aland@freeradius.org)
24 */
25RCSID("$Id: b45543b357714504d5ca0ceeab7a37f931a36fee $")
26
27#include <freeradius-devel/io/channel.h>
28#include <freeradius-devel/util/debug.h>
29
30#ifdef HAVE_STDATOMIC_H
31# include <stdatomic.h>
32#else
33# include <freeradius-devel/util/stdatomic.h>
34#endif
35
36/*
37 * Debugging, mainly for channel_test
38 */
39#ifdef DEBUG_CHANNEL
40#define MPRINT(...) fprintf(stdout, __VA_ARGS__)
41#else
42#define MPRINT(...)
43#endif
44
45/*
46 * We disable this until we fix all of the signaling issues...
47 */
48#define ENABLE_SKIPS (0)
49
54
55#ifdef DEBUG_CHANNEL
56static fr_table_num_sorted_t const channel_direction[] = {
57 { L("to responder"), TO_RESPONDER },
58 { L("to requestor"), TO_REQUESTOR },
59};
60size_t channel_direction_len = NUM_ELEMENTS(channel_direction);
61#endif
62
63#if 0
64#define SIGNAL_INTERVAL (1000000) //!< The minimum interval between responder signals.
65#endif
66
67/** Size of the atomic queues
68 *
69 * The queue reader MUST service the queue occasionally,
70 * otherwise the writer will not be able to write. If it's too
71 * low, the writer will fail. If it's too high, it will
72 * unnecessarily use memory. So we're better off putting it on
73 * the high side.
74 *
75 * The reader SHOULD service the queues at inter-packet latency.
76 * i.e. at 1M pps, the queue will get serviced every microsecond.
77 */
78#define ATOMIC_QUEUE_SIZE (1024)
79
94
95typedef struct {
96 fr_channel_signal_t signal; //!< the signal to send
97 uint64_t ack; //!< or the endpoint..
98 fr_channel_t *ch; //!< the channel
99 void *uctx; //!< arbitrary message data
101
102/** One end of a channel
103 *
104 * Consists of a kqueue descriptor, and an atomic queue.
105 * The atomic queue is there to get bulk data through, because it's more efficient
106 * than pushing 1M+ events per second through a kqueue.
107 */
108typedef struct {
109 fr_channel_direction_t direction; //!< Use for debug messages.
110
111 fr_control_t *control; //!< The control plane, consisting of an atomic queue and kqueue.
112
113 fr_ring_buffer_t *rb; //!< Ring buffer for control-plane messages.
114
115 void *uctx; //!< Worker context.
116
117 fr_channel_recv_callback_t recv; //!< callback for receiving messages
118 void *recv_uctx; //!< context for receiving messages
119
120 bool must_signal; //!< we need to signal the other end
121
122
123 uint64_t sequence; //!< Sequence number for this channel.
124 uint64_t ack; //!< Sequence number of the other end.
125 uint64_t their_view_of_my_sequence; //!< Should be clear.
126
127 uint64_t sequence_at_last_signal; //!< When we last signaled.
128
129 fr_atomic_queue_t *aq; //!< The queue of messages - visible only to this channel.
130
131 atomic_bool active; //!< Whether the channel is active.
132
133 fr_channel_stats_t stats; //!< channel statistics
135
137
138/** A full channel, which consists of two ends
139 *
140 * A channel consists of an I/O identifier that can be placed in kequeue
141 * and an atomic queue in each direction to allow for bidirectional communication.
142 */
144 fr_time_delta_t cpu_time; //!< Total time used by the responder for this channel.
145 fr_time_delta_t processing_time; //!< Time spent by the responder processing requests.
146
147 bool same_thread; //!< are both ends in the same thread?
148
149 fr_channel_end_t end[2]; //!< Two ends of the channel.
150};
151
153 { L("error"), FR_CHANNEL_ERROR },
154 { L("data-to-responder"), FR_CHANNEL_SIGNAL_DATA_TO_RESPONDER },
155 { L("data-to-requestor"), FR_CHANNEL_DATA_READY_REQUESTOR },
156 { L("open"), FR_CHANNEL_OPEN },
157 { L("close"), FR_CHANNEL_CLOSE },
158 { L("data-done-responder"), FR_CHANNEL_SIGNAL_DATA_DONE_RESPONDER },
159 { L("responder-sleeping"), FR_CHANNEL_SIGNAL_RESPONDER_SLEEPING },
160};
162
164 { L("high"), PRIORITY_HIGH },
165 { L("low"), PRIORITY_LOW },
166 { L("normal"), PRIORITY_NORMAL },
167 { L("now"), PRIORITY_NOW }
168};
170
171
172/** Create a new channel
173 *
174 * @param[in] ctx The talloc_ctx to allocate channel data in.
175 * @param[in] requestor control plane.
176 * @param[in] responder control plane.
177 * @param[in] same whether or not the channel is for the same thread
178 * @return
179 * - NULL on error
180 * - channel on success
181 */
182fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, fr_control_t *requestor, fr_control_t *responder, bool same)
183{
184 fr_time_t now;
185 fr_channel_t *ch;
186
187 ch = talloc_zero(ctx, fr_channel_t);
188 if (!ch) {
189 nomem:
190 fr_strerror_const("Failed allocating memory");
191 return NULL;
192 }
193
194 ch->same_thread = same;
195
198
200 if (!ch->end[TO_RESPONDER].aq) {
201 talloc_free(ch);
202 goto nomem;
203 }
204
206 if (!ch->end[TO_REQUESTOR].aq) {
207 talloc_free(ch);
208 goto nomem;
209 }
210
211 ch->end[TO_RESPONDER].control = responder;
212 ch->end[TO_REQUESTOR].control = requestor;
213
214 /*
215 * Create the ring buffer for the requestor to send
216 * control-plane messages to the responder, and vice-versa.
217 */
219 if (!ch->end[TO_RESPONDER].rb) {
220 rb_nomem:
221 fr_strerror_const_push("Failed allocating ring buffer");
222 talloc_free(ch);
223 return NULL;
224 }
225
227 if (!ch->end[TO_REQUESTOR].rb) {
228 talloc_free(ch);
229 goto rb_nomem;
230 }
231
232 /*
233 * Initialize all of the timers to now.
234 */
235 now = fr_time();
236
237 ch->end[TO_RESPONDER].stats.last_write = now;
240 atomic_store(&ch->end[TO_RESPONDER].active, true);
241
242 ch->end[TO_REQUESTOR].stats.last_write = now;
245 atomic_store(&ch->end[TO_REQUESTOR].active, true);
246
247 return ch;
248}
249
250
251/** Send a message via a kq user signal
252 *
253 * Note that the caller doesn't care about data in the event, that is
254 * sent via the atomic queue. The kevent code takes care of
255 * delivering the signal once, even if it's sent by multiple requestor
256 * threads.
257 *
258 * The thread watching the KQ knows which end it is. So when it gets
259 * the signal (and the channel pointer) it knows to look at end[0] or
260 * end[1]. We also send which end in 'which' (0, 1) to further help
261 * the recipient.
262 *
263 * @param[in] ch the channel.
264 * @param[in] when the data was ready. Typically taken from the message.
265 * @param[in] end of the channel that the message was written to.
266 * @param[in] which end of the channel (0/1).
267 * @return
268 * - <0 on error
269 * - 0 on success
270 */
272{
274
275 end->stats.last_sent_signal = when;
276 end->stats.signals++;
277 end->must_signal = false;
278
279 cc.signal = which;
280 cc.ack = end->ack;
281 cc.ch = ch;
282
283 MPRINT("Signalling %s, with %s\n",
284 fr_table_str_by_value(channel_direction, end->direction, "<INVALID>"),
285 fr_table_str_by_value(channel_signals, which, "<INVALID>"));
286
287 return fr_control_message_send(end->control, end->rb, FR_CONTROL_ID_CHANNEL, &cc, sizeof(cc));
288}
289
290#define IALPHA (8)
291#define RTT(_old, _new) fr_time_delta_wrap((fr_time_delta_unwrap(_new) + (fr_time_delta_unwrap(_old) * (IALPHA - 1))) / IALPHA)
292
293/** Send a request message into the channel
294 *
295 * The message should be initialized, other than "sequence" and "ack".
296 *
297 * This function automatically calls the recv_reply callback if there is a reply.
298 *
299 * @param[in] ch the channel to send the request on.
300 * @param[in] cd the message to send.
301 * @return
302 * - <0 on error
303 * - 0 on success
304 */
306{
307 uint64_t sequence;
308 fr_time_t when;
309 fr_time_delta_t message_interval;
310 fr_channel_end_t *requestor;
311
312 if (!fr_cond_assert_msg(atomic_load(&ch->end[TO_RESPONDER].active), "Channel not active")) return -1;
313
314 /*
315 * Same thread? Just call the "recv" function directly.
316 */
317 if (ch->same_thread) {
318 ch->end[TO_REQUESTOR].recv(ch, cd, ch->end[TO_REQUESTOR].recv_uctx);
319 return 0;
320 }
321
322 requestor = &(ch->end[TO_RESPONDER]);
323 when = cd->m.when;
324
325 sequence = requestor->sequence + 1;
326 cd->live.sequence = sequence;
327 cd->live.ack = requestor->ack;
328
329 /*
330 * Push the message onto the queue for the other end. If
331 * the push fails, the caller should try another queue.
332 */
333 if (!fr_atomic_queue_push(requestor->aq, cd)) {
334 fr_strerror_printf("Failed pushing to atomic queue - full. Queue contains %zu items",
335 fr_atomic_queue_size(requestor->aq));
336 while (fr_channel_recv_reply(ch));
337 return -1;
338 }
339
340 requestor->sequence = sequence;
341 message_interval = fr_time_sub(when, requestor->stats.last_write);
342
343 if (!fr_time_delta_ispos(requestor->stats.message_interval)) {
344 requestor->stats.message_interval = message_interval;
345 } else {
346 requestor->stats.message_interval = RTT(requestor->stats.message_interval, message_interval);
347 }
348
349 fr_assert_msg(fr_time_lteq(requestor->stats.last_write, when),
350 "Channel data timestamp (%" PRId64") older than last channel data sent (%" PRId64 ")",
351 fr_time_unwrap(when), fr_time_unwrap(requestor->stats.last_write));
352 requestor->stats.last_write = when;
353
354 requestor->stats.outstanding++;
355 requestor->stats.packets++;
356
357 MPRINT("REQUESTOR requests %"PRIu64", num_outstanding %"PRIu64"\n", requestor->stats.packets, requestor->stats.outstanding);
358
359#if ENABLE_SKIPS
360 /*
361 * We just sent the first packet. There can't possibly be a reply, so don't bother looking.
362 */
363 if (requestor->stats.outstanding == 1) {
364
365 /*
366 * There is at least one old packet which is
367 * outstanding, look for a reply.
368 */
369 } else if (requestor->stats.outstanding > 1) {
370 bool has_reply;
371
372 has_reply = fr_channel_recv_reply(ch);
373
374 if (has_reply) while (fr_channel_recv_reply(ch));
375
376 /*
377 * There's no reply yet, so we still have packets outstanding.
378 * Or, there is a reply, and there are more packets outstanding.
379 * Skip the signal.
380 */
381 if (!requestor->must_signal && (!has_reply || (has_reply && (requestor->stats.outstanding > 1)))) {
382 MPRINT("REQUESTOR SKIPS signal\n");
383 return 0;
384 }
385 }
386#endif
387
388 /*
389 * Tell the other end that there is new data ready.
390 *
391 * Ignore errors on signalling. The responder already has
392 * the packet in its inbound queue, so at some point, it
393 * will pick up the message.
394 */
395 MPRINT("REQUESTOR SIGNALS\n");
397 return 0;
398}
399
400/** Receive a reply message from the channel
401 *
402 * @param[in] ch the channel to read data from.
403 * @return
404 * - true if there was a message received
405 * - false if there are no more messages
406 */
408{
410 fr_channel_end_t *requestor;
412
413 fr_assert(ch->end[TO_RESPONDER].recv != NULL);
414
415 aq = ch->end[TO_REQUESTOR].aq;
416 requestor = &(ch->end[TO_RESPONDER]);
417
418 /*
419 * It's OK for the queue to be empty.
420 */
421 if (!fr_atomic_queue_pop(aq, (void **) &cd)) return false;
422
423 /*
424 * We want an exponential moving average for round trip
425 * time, where "alpha" is a number between [0,1)
426 *
427 * RTT_new = alpha * RTT_old + (1 - alpha) * RTT_sample
428 *
429 * BUT we use fixed-point arithmetic, so we need to use inverse alpha,
430 * which works out to the following equation:
431 *
432 * RTT_new = (RTT_sample + (ialpha - 1) * RTT_old) / ialpha
433 *
434 * NAKs have zero processing time, so we ignore them for
435 * the purpose of RTT.
436 */
437 if (fr_time_delta_ispos(cd->reply.processing_time)) {
438 ch->processing_time = RTT(ch->processing_time, cd->reply.processing_time);
439 }
440 ch->cpu_time = cd->reply.cpu_time;
441
442 /*
443 * Update the outbound channel with the knowledge that
444 * we've received one more reply, and with the responders
445 * ACK.
446 */
447 fr_assert(requestor->stats.outstanding > 0);
448 fr_assert(cd->live.sequence > requestor->ack);
449 fr_assert(cd->live.sequence <= requestor->sequence); /* must have fewer replies than requests */
450
451 requestor->stats.outstanding--;
452 requestor->ack = cd->live.sequence;
453 requestor->their_view_of_my_sequence = cd->live.ack;
454
456 requestor->stats.last_read_other = cd->m.when;
457
458 ch->end[TO_RESPONDER].recv(ch, cd, ch->end[TO_RESPONDER].recv_uctx);
459
460 return true;
461}
462
463
464/** Receive a request message from the channel
465 *
466 * @param[in] ch the channel
467 * @return
468 * - true if there was a message received
469 * - false if there are no more messages
470 */
472{
474 fr_channel_end_t *responder;
476
477 aq = ch->end[TO_RESPONDER].aq;
478 responder = &(ch->end[TO_REQUESTOR]);
479
480 /*
481 * It's OK for the queue to be empty.
482 */
483 if (!fr_atomic_queue_pop(aq, (void **) &cd)) return false;
484
485 fr_assert(cd->live.sequence > responder->ack);
486 fr_assert(cd->live.sequence >= responder->sequence); /* must have more requests than replies */
487
488 responder->stats.outstanding++;
489 responder->ack = cd->live.sequence;
490 responder->their_view_of_my_sequence = cd->live.ack;
491
493 responder->stats.last_read_other = cd->m.when;
494
495 ch->end[TO_REQUESTOR].recv(ch, cd, ch->end[TO_REQUESTOR].recv_uctx);
496
497 return true;
498}
499
500/** Send a reply message into the channel
501 *
502 * The message should be initialized, other than "sequence" and "ack".
503 *
504 * @param[in] ch the channel to send the reply on.
505 * @param[in] cd the message to send
506 * @return
507 * - <0 on error
508 * - 0 on success
509 */
511{
512 uint64_t sequence;
513 fr_time_t when;
514 fr_time_delta_t message_interval;
515 fr_channel_end_t *responder;
516
517 if (!fr_cond_assert_msg(atomic_load(&ch->end[TO_REQUESTOR].active), "Channel not active")) return -1;
518
519 /*
520 * Same thread? Just call the "recv" function directly.
521 */
522 if (ch->same_thread) {
523 ch->end[TO_RESPONDER].recv(ch, cd, ch->end[TO_RESPONDER].recv_uctx);
524 return 0;
525 }
526
527 responder = &(ch->end[TO_REQUESTOR]);
528
529 when = cd->m.when;
530
531 sequence = responder->sequence + 1;
532 cd->live.sequence = sequence;
533 cd->live.ack = responder->ack;
534
535 if (!fr_atomic_queue_push(responder->aq, cd)) {
536 fr_strerror_printf("Failed pushing to atomic queue - full. Queue contains %zu items",
537 fr_atomic_queue_size(responder->aq));
538 while (fr_channel_recv_request(ch));
539 return -1;
540 }
541
542 fr_assert(responder->stats.outstanding > 0);
543 responder->stats.outstanding--;
544 responder->stats.packets++;
545
546 MPRINT("\tRESPONDER replies %"PRIu64", num_outstanding %"PRIu64"\n", responder->stats.packets, responder->stats.outstanding);
547
548 responder->sequence = sequence;
549 message_interval = fr_time_sub(when, responder->stats.last_write);
550 if (!fr_time_delta_ispos(responder->stats.message_interval)) {
551 responder->stats.message_interval = message_interval;
552 } else {
553 responder->stats.message_interval = RTT(responder->stats.message_interval, message_interval);
554 }
555
556 fr_assert_msg(fr_time_lteq(responder->stats.last_write, when),
557 "Channel data timestamp (%" PRId64") older than last channel data sent (%" PRId64 ")",
558 fr_time_unwrap(when), fr_time_unwrap(responder->stats.last_write));
559 responder->stats.last_write = when;
560
561 /*
562 * Even if we think we have no more packets to process,
563 * the caller may have sent us one. Go check the input
564 * channel.
565 */
566 while (fr_channel_recv_request(ch));
567
568 /*
569 * No packets outstanding, we HAVE to signal the requestor
570 * thread.
571 */
572 if (responder->stats.outstanding == 0) {
574 return 0;
575 }
576
577 MPRINT("\twhen - last_read_other = %"PRIu64" - %"PRIu64" = %"PRIu64"\n", when, responder->stats.last_read_other, when - responder->stats.last_read_other);
578 MPRINT("\twhen - last signal = %"PRIu64" - %"PRIu64" = %"PRIu64"\n", when, responder->stats.last_sent_signal, when - responder->stats.last_sent_signal);
579 MPRINT("\tsequence - ack = %"PRIu64" - %"PRIu64" = %"PRIu64"\n", responder->sequence, responder->their_view_of_my_sequence, responder->sequence - responder->their_view_of_my_sequence);
580
581#ifdef __APPLE__
582 /*
583 * If we've sent them a signal since the last ACK, they
584 * will receive it, and process the packets. So we don't
585 * need to signal them again.
586 *
587 * But... this doesn't appear to work on the Linux
588 * libkqueue implementation.
589 */
590 if (responder->sequence_at_last_signal > responder->their_view_of_my_sequence) return 0;
591#endif
592
593 /*
594 * If we've received a new packet in the last while, OR
595 * we've sent a signal in the last while, then we don't
596 * need to send a new signal. But we DO send a signal if
597 * we haven't seen an ACK for a few packets.
598 *
599 * FIXME: make these limits configurable, or include
600 * predictions about packet processing time?
601 */
602 fr_assert(responder->their_view_of_my_sequence <= responder->sequence);
603#if 0
604 if (((responder->sequence - their_view_of_my_sequence) <= 1000) &&
605 ((when - responder->stats.last_read_other < SIGNAL_INTERVAL) ||
606 ((when - responder->stats.last_sent_signal) < SIGNAL_INTERVAL))) {
607 MPRINT("\tRESPONDER SKIPS signal\n");
608 return 0;
609 }
610#endif
611
612 MPRINT("\tRESPONDER SIGNALS num_outstanding %"PRIu64"\n", responder->stats.outstanding);
614 return 0;
615}
616
617
618/** Don't send a reply message into the channel
619 *
620 * The message should be the one we received from the network.
621 *
622 * @param[in] ch the channel on which we're dropping a packet
623 * @return
624 * - <0 on error
625 * - 0 on success
626 */
628{
629 fr_channel_end_t *responder;
630
631 responder = &(ch->end[TO_REQUESTOR]);
632
633 responder->sequence++;
634 return 0;
635}
636
637
638
639/** Signal a channel that the responder is sleeping
640 *
641 * This function should be called from the responders idle loop.
642 * i.e. only when it has nothing else to do.
643 *
644 * @param[in] ch the channel to signal we're no longer listening on.
645 * @return
646 * - <0 on error
647 * - 0 on success
648 */
650{
651 fr_channel_end_t *responder;
653
654 responder = &(ch->end[TO_REQUESTOR]);
655
656 /*
657 * We don't have any outstanding requests to process for
658 * this channel, don't signal the network thread that
659 * we're sleeping. It already knows.
660 */
661 if (responder->stats.outstanding == 0) return 0;
662
663 responder->stats.signals++;
664
666 cc.ack = responder->ack;
667 cc.ch = ch;
668
669 MPRINT("\tRESPONDER SLEEPING num_outstanding %"PRIu64", packets in %"PRIu64", packets out %"PRIu64"\n", responder->stats.outstanding,
670 ch->end[TO_RESPONDER].stats.packets, responder->stats.packets);
671 return fr_control_message_send(responder->control, responder->rb, FR_CONTROL_ID_CHANNEL, &cc, sizeof(cc));
672}
673
674
675/** Service a control-plane message
676 *
677 * @param[in] when The current time.
678 * @param[out] p_channel The channel which should be serviced.
679 * @param[out] uctx_out The uctx sent in the control message.
680 * @param[in] data The control message.
681 * @param[in] data_size The size of the control message.
682 * @return
683 * - FR_CHANNEL_ERROR on error
684 * - FR_CHANNEL_NOOP, on do nothing
685 * - FR_CHANNEL_DATA_READY on data ready
686 * - FR_CHANNEL_OPEN when a channel has been opened and sent to us
687 * - FR_CHANNEL_CLOSE when a channel should be closed
688 */
689fr_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)
690{
691 int rcode;
692#if ENABLE_SKIPS
693 uint64_t ack;
694#endif
698 fr_channel_end_t *requestor;
699 fr_channel_t *ch;
700
701 fr_assert(data_size == sizeof(cc));
702 memcpy(&cc, data, data_size);
703
704 cs = cc.signal;
705#if ENABLE_SKIPS
706 ack = cc.ack;
707#endif
708 *p_channel = ch = cc.ch;
709 if (uctx_out) *uctx_out = cc.uctx;
710
711 switch (cs) {
712 /*
713 * These all have the same numbers as the channel
714 * events, and have no extra processing. We just
715 * return them as-is.
716 */
722 MPRINT("channel got %d\n", cs);
723 return (fr_channel_event_t) cs;
724
725 /*
726 * Only sent by the responder. Both of these
727 * situations are largely the same, except for
728 * return codes.
729 */
731 MPRINT("channel got data_done_responder\n");
733 ch->end[TO_RESPONDER].must_signal = true;
734 break;
735
737 MPRINT("channel got responder_sleeping\n");
738 ce = FR_CHANNEL_NOOP;
739 ch->end[TO_RESPONDER].must_signal = true;
740 break;
741 }
742
743 /*
744 * Compare their ACK to the last sequence we
745 * sent. If it's different, we signal the responder
746 * to wake up.
747 */
748 requestor = &ch->end[TO_RESPONDER];
749
750 /*
751 * We have told this end to close, so there is nothing left to
752 * wake it for: fr_channel_send_request() refuses on an inactive
753 * end, so no more data can be queued. Signalling it anyway
754 * writes into the control plane which the responder may have
755 * already freed (this has been observed and caused a crash-on-exit).
756 *
757 * Servicing the message is still correct: it was in flight
758 * before the close, and we must keep servicing to receive the
759 * ack at all. The load is relaxed because only this thread
760 * clears "active" for this end.
761 */
762 if (unlikely(!atomic_load_explicit(&requestor->active, memory_order_relaxed))) return ce;
763
764#if ENABLE_SKIPS
765 if (!requestor->must_signal && (ack == requestor->sequence)) {
766 MPRINT("REQUESTOR SKIPS signal AFTER CE %d num_outstanding %"PRIu64"\n", cs, requestor->stats.outstanding);
767 MPRINT("REQUESTOR has ack %"PRIu64", my seq %"PRIu64" my_view %"PRIu64"\n", ack, requestor->sequence, requestor->their_view_of_my_sequence);
768 return ce;
769 }
770
771 /*
772 * The responder is sleeping or done. There are more
773 * packets available, so we signal it to wake up again.
774 */
775 fr_assert(ack <= requestor->sequence);
776#endif
777
778 /*
779 * We're signaling it again...
780 */
781 requestor->stats.resignals++;
782
783 /*
784 * The responder hasn't seen our last few packets. Signal
785 * that there is data ready.
786 */
787 MPRINT("REQUESTOR SIGNALS AFTER CE %d\n", cs);
789 if (rcode < 0) return FR_CHANNEL_ERROR;
790
791 return ce;
792}
793
794
795/** Service a control-plane event.
796 *
797 * The channels use control planes for internal signaling. Note that
798 * the caller does NOT pass the channel into this function. Instead,
799 * the channel is taken from the kevent.
800 *
801 * @param[in] ch The channel to service.
802 * @param[in] c The control plane on which we received the kev.
803 * @param[in] kev The kevent data, should get passed to the control plane.
804 * @return
805 * - <0 on error
806 * - 0 on success
807 */
808int fr_channel_service_kevent(fr_channel_t *ch, fr_control_t *c, UNUSED struct kevent const *kev)
809{
810 (void) talloc_get_type_abort(ch, fr_channel_t);
811
812 if (c == ch->end[TO_RESPONDER].control) {
814 } else {
816 }
817
818 return 0;
819}
820
821
822/** Check if a channel is active.
823 *
824 * A channel may be closed by either end. If so, it stays alive (but
825 * inactive) until both ends acknowledge the close.
826 *
827 * @param[in] ch the channel
828 * @return
829 * - false the channel is closing.
830 * - true the channel is active
831 */
836
837/** Signal a responder that the channel is closing
838 *
839 * @param[in] ch The channel.
840 * @return
841 * - <0 on error
842 * - 0 on success
843 */
845{
846 int ret;
847 bool active;
849
850 active = atomic_load(&ch->end[TO_RESPONDER].active);
851 if (!active) return 0; /* Already signalled to close */
852
853 atomic_store(&ch->end[TO_RESPONDER].active, false); /* Prevent further requests */
854
855 (void) talloc_get_type_abort(ch, fr_channel_t);
856
858 cc.ack = TO_RESPONDER;
859 cc.ch = ch;
860
862 ch->end[TO_RESPONDER].rb, FR_CONTROL_ID_CHANNEL, &cc, sizeof(cc));
863
864 return ret;
865}
866
867/** Discard any requests the requestor queued but we never received
868 *
869 * The messages belong to the requestor's message set, and it cannot reclaim them
870 * until they are marked done. A responder that closes with requests still in
871 * the queue would strand them there, and with them the memory they were
872 * allocated from.
873 *
874 * @param[in] ch to discard the queued requests of.
875 * @return the number of requests discarded.
876 */
878{
881 unsigned int num = 0;
882
883 while (fr_atomic_queue_pop(aq, (void **) &cd)) {
884 fr_message_done(&cd->m);
885 num++;
886 }
887
888 return num;
889}
890
891/** Acknowledge that the channel is closing
892 *
893 * @param[in] ch The channel.
894 * @return
895 * - <0 on error
896 * - 0 on success
897 */
899{
900 int ret;
901 bool active;
902
904
905 active = atomic_load(&ch->end[TO_REQUESTOR].active);
906 if (!active) return 0; /* Already signalled to close */
907
908 atomic_store(&ch->end[TO_REQUESTOR].active, false); /* Prevent further responses */
909
910 (void) talloc_get_type_abort(ch, fr_channel_t);
911
913 cc.ack = TO_REQUESTOR;
914 cc.ch = ch;
915
917 ch->end[TO_REQUESTOR].rb, FR_CONTROL_ID_CHANNEL, &cc, sizeof(cc));
918
919 return ret;
920}
921
922/** Add responder-specific data to a channel
923 *
924 * @param[in] ch The channel.
925 * @param[in] uctx The context to add.
926 */
928{
929 (void) talloc_get_type_abort(ch, fr_channel_t);
930
931 ch->end[TO_REQUESTOR].uctx = uctx;
932}
933
934
935/** Get responder-specific data from a channel
936 *
937 * @param[in] ch The channel.
938 */
940{
941 (void) talloc_get_type_abort(ch, fr_channel_t);
942
943 return ch->end[TO_REQUESTOR].uctx;
944}
945
946
947/** Add network-specific data to a channel
948 *
949 * @param[in] ch The channel.
950 * @param[in] uctx The context to add.
951 */
953{
954 (void) talloc_get_type_abort(ch, fr_channel_t);
955
956 ch->end[TO_RESPONDER].uctx = uctx;
957}
958
959
960/** Get network-specific data from a channel
961 *
962 * @param[in] ch The channel.
963 */
965{
966 (void) talloc_get_type_abort(ch, fr_channel_t);
967
968 return ch->end[TO_RESPONDER].uctx;
969}
970
971
973{
974 ch->end[TO_RESPONDER].recv = recv_reply;
975 ch->end[TO_RESPONDER].recv_uctx = uctx;
976
977 return 0;
978}
979
981{
982 ch->end[TO_REQUESTOR].recv = recv_request;
983 ch->end[TO_REQUESTOR].recv_uctx = uctx;
984 return 0;
985}
986
987/** Send a channel to a responder
988 *
989 * @param[in] ch The channel.
990 * @return
991 * - <0 on error
992 * - 0 on success
993 */
995{
997
999 cc.ack = 0;
1000 cc.ch = ch;
1001 cc.uctx = uctx;
1002
1004}
1005
1006void fr_channel_stats_log(fr_channel_t const *ch, fr_log_t const *log, char const *file, int line)
1007{
1008 fr_log(log, L_INFO, file, line, "requestor\n");
1009 fr_log(log, L_INFO, file, line, "\tsignals sent = %" PRIu64 "\n", ch->end[TO_RESPONDER].stats.signals);
1010 fr_log(log, L_INFO, file, line, "\tsignals re-sent = %" PRIu64 "\n", ch->end[TO_RESPONDER].stats.resignals);
1011 fr_log(log, L_INFO, file, line, "\tkevents checked = %" PRIu64 "\n", ch->end[TO_RESPONDER].stats.kevents);
1012 fr_log(log, L_INFO, file, line, "\toutstanding = %" PRIu64 "\n", ch->end[TO_RESPONDER].stats.outstanding);
1013 fr_log(log, L_INFO, file, line, "\tpackets processed = %" PRIu64 "\n", ch->end[TO_RESPONDER].stats.packets);
1014 fr_log(log, L_INFO, file, line, "\tmessage interval (RTT) = %" PRIu64 "\n", fr_time_delta_unwrap(ch->end[TO_RESPONDER].stats.message_interval));
1015 fr_log(log, L_INFO, file, line, "\tlast write = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_RESPONDER].stats.last_write));
1016 fr_log(log, L_INFO, file, line, "\tlast read other end = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_RESPONDER].stats.last_read_other));
1017 fr_log(log, L_INFO, file, line, "\tlast signal other = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_RESPONDER].stats.last_sent_signal));
1018
1019 fr_log(log, L_INFO, file, line, "responder\n");
1020 fr_log(log, L_INFO, file, line, "\tsignals sent = %" PRIu64"\n", ch->end[TO_REQUESTOR].stats.signals);
1021 fr_log(log, L_INFO, file, line, "\tkevents checked = %" PRIu64 "\n", ch->end[TO_REQUESTOR].stats.kevents);
1022 fr_log(log, L_INFO, file, line, "\tpackets processed = %" PRIu64 "\n", ch->end[TO_REQUESTOR].stats.packets);
1023 fr_log(log, L_INFO, file, line, "\tmessage interval (RTT) = %" PRIu64 "\n", fr_time_delta_unwrap(ch->end[TO_REQUESTOR].stats.message_interval));
1024 fr_log(log, L_INFO, file, line, "\tlast write = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_REQUESTOR].stats.last_write));
1025 fr_log(log, L_INFO, file, line, "\tlast read other end = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_REQUESTOR].stats.last_read_other));
1026 fr_log(log, L_INFO, file, line, "\tlast signal other = %" PRIu64 "\n", fr_time_unwrap(ch->end[TO_REQUESTOR].stats.last_sent_signal));
1027}
int const char * file
Definition acutest.h:702
int const char int line
Definition acutest.h:702
bool fr_atomic_queue_pop(fr_atomic_queue_t *aq, void **p_data)
Pop a pointer from the atomic queue.
size_t fr_atomic_queue_size(fr_atomic_queue_t *aq)
bool fr_atomic_queue_push(fr_atomic_queue_t *aq, void *data)
Push a pointer into the atomic queue.
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 L(_str)
Helper for initialising arrays of string literals.
Definition build.h:228
#define unlikely(_x)
Definition build.h:455
#define UNUSED
Definition build.h:384
#define NUM_ELEMENTS(_t)
Definition build.h:406
fr_atomic_queue_t * aq
The queue of messages - visible only to this channel.
Definition channel.c:129
atomic_bool active
Whether the channel is active.
Definition channel.c:131
#define MPRINT(...)
Definition channel.c:42
void * fr_channel_requestor_uctx_get(fr_channel_t *ch)
Get network-specific data from a channel.
Definition channel.c:964
fr_channel_signal_t
Definition channel.c:80
@ FR_CHANNEL_SIGNAL_DATA_DONE_RESPONDER
Definition channel.c:91
@ FR_CHANNEL_SIGNAL_DATA_TO_REQUESTOR
Definition channel.c:83
@ FR_CHANNEL_SIGNAL_DATA_TO_RESPONDER
Definition channel.c:82
@ FR_CHANNEL_SIGNAL_RESPONDER_SLEEPING
Definition channel.c:92
@ FR_CHANNEL_SIGNAL_ERROR
Definition channel.c:81
@ FR_CHANNEL_SIGNAL_OPEN
Definition channel.c:84
@ FR_CHANNEL_SIGNAL_CLOSE
Definition channel.c:85
uint64_t sequence_at_last_signal
When we last signaled.
Definition channel.c:127
uint64_t sequence
Sequence number for this channel.
Definition channel.c:123
bool must_signal
we need to signal the other end
Definition channel.c:120
fr_channel_signal_t signal
the signal to send
Definition channel.c:96
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
fr_channel_direction_t direction
Use for debug messages.
Definition channel.c:109
unsigned int fr_channel_responder_discard(fr_channel_t *ch)
Discard any requests the requestor queued but we never received.
Definition channel.c:877
size_t channel_signals_len
Definition channel.c:161
#define RTT(_old, _new)
Definition channel.c:291
void * uctx
Worker context.
Definition channel.c:115
size_t channel_packet_priority_len
Definition channel.c:169
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
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:182
fr_channel_direction_t
Definition channel.c:50
@ TO_RESPONDER
Definition channel.c:51
@ TO_REQUESTOR
Definition channel.c:52
fr_table_num_sorted_t const channel_packet_priority[]
Definition channel.c:163
static int fr_channel_data_ready(fr_channel_t *ch, fr_time_t when, fr_channel_end_t *end, fr_channel_signal_t which)
Send a message via a kq user signal.
Definition channel.c:271
fr_ring_buffer_t * rb
Ring buffer for control-plane messages.
Definition channel.c:113
#define ATOMIC_QUEUE_SIZE
Size of the atomic queues.
Definition channel.c:78
fr_channel_stats_t stats
channel statistics
Definition channel.c:133
fr_control_t * control
The control plane, consisting of an atomic queue and kqueue.
Definition channel.c:111
int fr_channel_signal_open(fr_channel_t *ch, void *uctx)
Send a channel to a responder.
Definition channel.c:994
bool same_thread
are both ends in the same thread?
Definition channel.c:147
int fr_channel_set_recv_reply(fr_channel_t *ch, fr_channel_recv_callback_t recv_reply, void *uctx)
Definition channel.c:972
uint64_t ack
or the endpoint..
Definition channel.c:97
void * fr_channel_responder_uctx_get(fr_channel_t *ch)
Get responder-specific data from a channel.
Definition channel.c:939
bool fr_channel_recv_request(fr_channel_t *ch)
Receive a request message from the channel.
Definition channel.c:471
int fr_channel_null_reply(fr_channel_t *ch)
Don't send a reply message into the channel.
Definition channel.c:627
void fr_channel_requestor_uctx_add(fr_channel_t *ch, void *uctx)
Add network-specific data to a channel.
Definition channel.c:952
fr_time_delta_t cpu_time
Total time used by the responder for this channel.
Definition channel.c:144
int fr_channel_service_kevent(fr_channel_t *ch, fr_control_t *c, UNUSED struct kevent const *kev)
Service a control-plane event.
Definition channel.c:808
void fr_channel_responder_uctx_add(fr_channel_t *ch, void *uctx)
Add responder-specific data to a channel.
Definition channel.c:927
int fr_channel_responder_sleeping(fr_channel_t *ch)
Signal a channel that the responder is sleeping.
Definition channel.c:649
int fr_channel_set_recv_request(fr_channel_t *ch, fr_channel_recv_callback_t recv_request, void *uctx)
Definition channel.c:980
int fr_channel_send_reply(fr_channel_t *ch, fr_channel_data_t *cd)
Send a reply message into the channel.
Definition channel.c:510
fr_channel_end_t end[2]
Two ends of the channel.
Definition channel.c:149
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_t * ch
the channel
Definition channel.c:98
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
uint64_t their_view_of_my_sequence
Should be clear.
Definition channel.c:125
bool fr_channel_active(fr_channel_t *ch)
Check if a channel is active.
Definition channel.c:832
fr_time_delta_t processing_time
Time spent by the responder processing requests.
Definition channel.c:145
void * uctx
arbitrary message data
Definition channel.c:99
int fr_channel_responder_ack_close(fr_channel_t *ch)
Acknowledge that the channel is closing.
Definition channel.c:898
uint64_t ack
Sequence number of the other end.
Definition channel.c:124
fr_channel_recv_callback_t recv
callback for receiving messages
Definition channel.c:117
void * recv_uctx
context for receiving messages
Definition channel.c:118
One end of a channel.
Definition channel.c:108
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_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
fr_time_delta_t message_interval
Interval between messages.
Definition channel.h:94
#define FR_CONTROL_ID_CHANNEL
Definition channel.h:67
#define PRIORITY_NORMAL
Definition channel.h:153
uint64_t packets
Number of actual data packets.
Definition channel.h:88
uint64_t resignals
Number of signals resent.
Definition channel.h:86
void(* fr_channel_recv_callback_t)(fr_channel_t *ch, fr_channel_data_t *cd, void *uctx)
Definition channel.h:171
uint64_t outstanding
Number of outstanding requests with no reply.
Definition channel.h:84
fr_time_t last_sent_signal
The last time when we signaled the other end.
Definition channel.h:96
#define PRIORITY_NOW
Definition channel.h:151
fr_time_t last_read_other
Last time we successfully read a message from the other the channel.
Definition channel.h:93
#define PRIORITY_HIGH
Definition channel.h:152
fr_time_t last_write
Last write to the channel.
Definition channel.h:92
uint64_t kevents
Number of times we've looked at kevents.
Definition channel.h:90
uint64_t signals
Number of kevent signals we've sent.
Definition channel.h:85
#define PRIORITY_LOW
Definition channel.h:154
Channel information which is added to a message.
Definition channel.h:106
Statistics for the channel.
Definition channel.h:83
#define FR_CONTROL_MAX_SIZE
Definition control.h:51
#define FR_CONTROL_MAX_MESSAGES
Definition control.h:50
static fr_atomic_queue_t ** aq
#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
talloc_free(hp)
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
The control structure.
Definition control.c:76
#define fr_time()
Definition event.c:60
#define fr_log(_log, _lvl, _file, _line, _fmt,...)
Definition log.h:172
@ L_INFO
Informational message.
Definition log.h:52
int fr_message_done(fr_message_t *m)
Mark a message as done.
Definition message.c:195
fr_time_t when
when this message was sent
Definition message.h:47
#define fr_assert(_expr)
Definition rad_assert.h:37
fr_ring_buffer_t * fr_ring_buffer_create(TALLOC_CTX *ctx, size_t size)
Create a ring buffer.
Definition ring_buffer.c:64
#define atomic_store(object, desired)
Definition stdatomic.h:345
@ memory_order_relaxed
Definition stdatomic.h:127
#define atomic_load_explicit(object, order)
Definition stdatomic.h:312
#define atomic_load(object)
Definition stdatomic.h:343
Definition log.h:93
#define fr_table_str_by_value(_table, _number, _def)
Convert an integer to a string.
Definition table.h:804
An element in a lexicographically sorted array of name to num mappings.
Definition table.h:49
static int64_t fr_time_delta_unwrap(fr_time_delta_t time)
Definition time.h:154
static int64_t fr_time_unwrap(fr_time_t time)
Definition time.h:146
#define fr_time_lteq(_a, _b)
Definition time.h:240
#define fr_time_delta_ispos(_a)
Definition time.h:290
#define fr_time_sub(_a, _b)
Subtract one time from another.
Definition time.h:229
A time delta, a difference in time measured in nanoseconds.
Definition time.h:80
"server local" time.
Definition time.h:69
#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