From: Alan T. DeKok Date: Mon, 12 Dec 2016 19:30:14 +0000 (-0500) Subject: move to control-plane signaling for channels X-Git-Url: http://git.ipfire.org/cgi-bin/gitweb.cgi?a=commitdiff_plain;h=0365b345b3946527319e47fefa9722e6df705bf7;p=thirdparty%2Ffreeradius-server.git move to control-plane signaling for channels --- diff --git a/src/tests/util/channel_test.c b/src/tests/util/channel_test.c index 05ac568c6de..decc59be2c5 100644 --- a/src/tests/util/channel_test.c +++ b/src/tests/util/channel_test.c @@ -206,11 +206,13 @@ check_close: * Service the events. */ for (i = 0; i < num_events; i++) { - (void) fr_channel_service_kevent(aq_master, &events[i]); + (void) fr_channel_service_kevent(channel, aq_master, &events[i]); } now = fr_time(); + MPRINT1("Master servicing control-plane aq %p\n", aq_master); + while ((ce = fr_channel_service_aq(aq_master, now, &new_channel)) != FR_CHANNEL_EMPTY) { MPRINT1("Master got channel event %d\n", ce); @@ -320,11 +322,13 @@ static void *channel_worker(void *arg) if (num_events == 0) continue; for (i = 0; i < num_events; i++) { - (void) fr_channel_service_kevent(aq_worker, &events[i]); + (void) fr_channel_service_kevent(channel, aq_worker, &events[i]); } now = fr_time(); + MPRINT1("\tWorker servicing control-plane aq %p\n", aq_worker); + while ((ce = fr_channel_service_aq(aq_worker, now, &new_channel)) != FR_CHANNEL_EMPTY) { fr_channel_data_t *cd, *reply; @@ -405,11 +409,12 @@ static void *channel_worker(void *arg) break; case FR_CHANNEL_NOOP: + MPRINT1("\tWorker got NOOP\n"); rad_assert(new_channel == channel); break; default: - fprintf(stderr, "Worker got unexpected CE %d\n", ce); + fprintf(stderr, "\tWorker got unexpected CE %d\n", ce); /* * Not written yet! diff --git a/src/util/channel.c b/src/util/channel.c index 26ed6b318dc..27f352870ce 100644 --- a/src/util/channel.c +++ b/src/util/channel.c @@ -25,6 +25,7 @@ RCSID("$Id$") #include +#include #include /* @@ -93,7 +94,9 @@ typedef struct fr_channel_control_t { typedef struct fr_channel_end_t { int kq; //!< the kqueue associated with the channel - fr_atomic_queue_t *aq_control; //!< the control channel queue + fr_atomic_queue_t *aq_control; //!< the control plane queue - global to the thread + + fr_control_t *control; //!< the control plane int num_outstanding; //!< number of outstanding requests with no reply @@ -114,7 +117,7 @@ typedef struct fr_channel_end_t { fr_time_t last_sent_signal; //!< the last time when we signaled the other end - fr_atomic_queue_t *aq; //!< the queue of messages + fr_atomic_queue_t *aq; //!< the queue of messages - visible only to this channel } fr_channel_end_t; /** @@ -130,43 +133,6 @@ typedef struct fr_channel_t { } fr_channel_t; -static int fr_channel_add_kevent_worker(struct kevent *kev, int size) -{ - if (size < 3) return -1; - - EV_SET(&kev[0], FR_CHANNEL_SIGNAL_OPEN, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - EV_SET(&kev[1], FR_CHANNEL_SIGNAL_CLOSE, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - EV_SET(&kev[2], FR_CHANNEL_SIGNAL_DATA_TO_WORKER, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - - return 3; -} - - -static int fr_channel_add_kevent_receiver(struct kevent *kev, int size) -{ - if (size < 4) return -1; - - EV_SET(&kev[0], FR_CHANNEL_SIGNAL_WORKER_SLEEPING, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - EV_SET(&kev[1], FR_CHANNEL_SIGNAL_CLOSE, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - EV_SET(&kev[2], FR_CHANNEL_SIGNAL_DATA_FROM_WORKER, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - EV_SET(&kev[3], FR_CHANNEL_SIGNAL_DATA_DONE_WORKER, EVFILT_USER, EV_FLAG, NOTE_FFNOP, 0, NULL); - - return 4; -} - - -static int fr_channel_kevent_signal(fr_channel_end_t *end, fr_channel_control_t *cc) -{ - struct kevent kev; - - EV_SET(&kev, cc->signal, EVFILT_USER, 0, NOTE_TRIGGER | NOTE_FFCOPY, cc->ack, cc->ch); - - end->sequence_at_last_signal = end->sequence; - - return kevent(end->kq, &kev, 1, NULL, 0, NULL); -} - - /** Create a new channel * * @param[in] ctx the talloc_ctx for the channel @@ -183,8 +149,6 @@ fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_ { fr_time_t when; fr_channel_t *ch; - int num_events; - struct kevent events[4]; ch = talloc_zero(ctx, fr_channel_t); if (!ch) return NULL; @@ -219,19 +183,19 @@ fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_ ch->end[FROM_WORKER].last_read_other = when; ch->end[FROM_WORKER].last_sent_signal = when; - ch->active = true; - - num_events = fr_channel_add_kevent_worker(events, 4); - if (kevent(kq_worker, events, num_events, NULL, 0, NULL) < 0) { + ch->end[TO_WORKER].control = fr_control_create(ctx, + ch->end[TO_WORKER].kq, + ch->end[TO_WORKER].aq_control); + if (!ch->end[TO_WORKER].control) { talloc_free(ch); return NULL; } - num_events = fr_channel_add_kevent_receiver(events, 4); - if (kevent(kq_master, events, num_events, NULL, 0, NULL) < 0) { - talloc_free(ch); - return NULL; - } + MPRINT("Master CONTROL aq_master %p aq_worker %p\n", aq_master, aq_worker); + + MPRINT("Master CONTROL %p aq %p\n", ch->end[TO_WORKER].control, ch->end[TO_WORKER].aq_control); + + ch->active = true; return ch; } @@ -268,7 +232,7 @@ static int fr_channel_data_ready(fr_channel_t *ch, fr_time_t when, fr_channel_en cc.ack = end->ack; cc.ch = ch; - return fr_channel_kevent_signal(end, &cc); + return fr_control_message_send(end->control, &cc, sizeof(cc)); } #define IALPHA (8) @@ -563,7 +527,7 @@ int fr_channel_worker_sleeping(fr_channel_t *ch) cc.ack = end->ack; cc.ch = ch; - return fr_channel_kevent_signal(end, &cc); + return fr_control_message_send(end->control, &cc, sizeof(cc)); } @@ -585,21 +549,21 @@ fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t when, { int rcode; uint64_t ack; - fr_channel_control_t *cc; + ssize_t data_size; + fr_channel_control_t cc; fr_channel_signal_t cs; fr_channel_event_t ce = FR_CHANNEL_ERROR; fr_channel_end_t *end; fr_channel_t *ch; - if (!fr_atomic_queue_pop(aq, (void **) &cc)) { - *p_channel = NULL; - return FR_CHANNEL_EMPTY; - } + data_size = fr_control_message_pop(aq, &cc, sizeof(cc)); + if (data_size == 0) return FR_CHANNEL_EMPTY; + + rad_assert(data_size == sizeof(cc)); - cs = cc->signal; - ack = cc->ack; - *p_channel = ch = cc->ch; - talloc_free(cc); + cs = cc.signal; + ack = cc.ack; + *p_channel = ch = cc.ch; switch (cs) { /* @@ -669,40 +633,29 @@ fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t when, * event. Note that the caller does NOT pass the channel into this * function. Instead, the channel is taken from the kevent. * + * @param[in] ch the channel to service * @param[in] aq the atomic queue on which we receive control-plane messages * @param[in] kev the event of type EVFILT_USER * @return * - <0 on error * - 0 on success */ -int fr_channel_service_kevent(fr_atomic_queue_t *aq, struct kevent const *kev) +int fr_channel_service_kevent(fr_channel_t *ch, fr_atomic_queue_t *aq, UNUSED struct kevent const *kev) { - fr_channel_t *ch; - fr_channel_control_t *cc; - - rad_assert(kev->filter == EVFILT_USER); - - cc = talloc(aq, fr_channel_control_t); - rad_assert(cc != NULL); - - cc->signal = kev->ident; - cc->ack = (uint64_t) kev->data; - cc->ch = ch = kev->udata; - #ifndef NDEBUG talloc_get_type_abort(ch, fr_channel_t); #endif - + + if (fr_control_message_service_kevent(aq, kev) == 0) { + return 0; + } + if (aq == ch->end[TO_WORKER].aq_control) { ch->end[TO_WORKER].num_kevents++; } else { ch->end[FROM_WORKER].num_kevents++; } - if (!fr_atomic_queue_push(aq, cc)) { - return -1; - } - return 0; } @@ -739,7 +692,7 @@ int fr_channel_signal_worker_close(fr_channel_t *ch) cc.ack = TO_WORKER; cc.ch = ch; - return fr_channel_kevent_signal(&ch->end[TO_WORKER], &cc); + return fr_control_message_send(ch->end[TO_WORKER].control, &cc, sizeof(cc)); } /** Acknowledge that the channel is closing @@ -759,7 +712,7 @@ int fr_channel_ack_worker_close(fr_channel_t *ch) cc.ack = FROM_WORKER; cc.ch = ch; - return fr_channel_kevent_signal(&ch->end[FROM_WORKER], &cc); + return fr_control_message_send(ch->end[FROM_WORKER].control, &cc, sizeof(cc)); } /** Send a channel to a KQ @@ -777,7 +730,7 @@ int fr_channel_signal_open(fr_channel_t *ch) cc.ack = 0; cc.ch = ch; - return fr_channel_kevent_signal(&ch->end[TO_WORKER], &cc); + return fr_control_message_send(ch->end[TO_WORKER].control, &cc, sizeof(cc)); } void fr_channel_debug(fr_channel_t *ch, FILE *fp) @@ -806,9 +759,8 @@ void fr_channel_debug(fr_channel_t *ch, FILE *fp) * - <0 on error * - 0 on success */ -int fr_channel_worker_receive_open(UNUSED TALLOC_CTX *ctx, UNUSED fr_channel_t *ch) +int fr_channel_worker_receive_open(TALLOC_CTX *ctx, fr_channel_t *ch) { -#if 0 #ifndef NDEBUG talloc_get_type_abort(ch, fr_channel_t); #endif @@ -821,7 +773,8 @@ int fr_channel_worker_receive_open(UNUSED TALLOC_CTX *ctx, UNUSED fr_channel_t * if (!ch->end[FROM_WORKER].control) { return -1; } -#endif + + MPRINT("\tWorker CONTROL %p\n", ch->end[FROM_WORKER].control); return 0; } diff --git a/src/util/channel.h b/src/util/channel.h index bbfa8be6a00..5f926e150e2 100644 --- a/src/util/channel.h +++ b/src/util/channel.h @@ -115,7 +115,7 @@ fr_channel_data_t *fr_channel_recv_reply(fr_channel_t *ch) CC_HINT(nonnull); int fr_channel_worker_sleeping(fr_channel_t *ch) CC_HINT(nonnull); -int fr_channel_service_kevent(fr_atomic_queue_t *aq, struct kevent const *kev) CC_HINT(nonnull); +int fr_channel_service_kevent(fr_channel_t *ch, fr_atomic_queue_t *aq, struct kevent const *kev) CC_HINT(nonnull); fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t when, fr_channel_t **p_channel) CC_HINT(nonnull); bool fr_channel_active(fr_channel_t *ch) CC_HINT(nonnull); diff --git a/src/util/control.c b/src/util/control.c index 81caad426f9..395975fdb59 100644 --- a/src/util/control.c +++ b/src/util/control.c @@ -342,22 +342,20 @@ ssize_t fr_control_message_pop(fr_atomic_queue_t *aq, void *data, size_t data_si } -/** Receive a control-plane message +/** Service a control-plane kevent * * This function is called ONLY from the receiving thread. * * @param[in] aq the recipients atomic queue for control-plane messages * @param[in] kev the kevent for this receiver - * @param[in,out] data where the data is stored - * @param[in] data_size the size of the buffer where we store the data. * @return - * - <0 the size of the data we need to read the next message + * - <0 error * - 0 this kevent is not for us. - * - >0 the amount of data we've read + * - >0 this kevent is for us */ -ssize_t fr_control_message_receive(fr_atomic_queue_t *aq, struct kevent const *kev, void *data, size_t data_size) +int fr_control_message_service_kevent(UNUSED fr_atomic_queue_t *aq, struct kevent const *kev) { if (kev->ident != FR_CONTROL_SIGNAL) return 0; - return fr_control_message_pop(aq, data, data_size); + return 0; } diff --git a/src/util/control.h b/src/util/control.h index 47897a36cea..341f7c755cd 100644 --- a/src/util/control.h +++ b/src/util/control.h @@ -47,7 +47,7 @@ void fr_control_free(fr_control_t *c); int fr_control_gc(fr_control_t *c) CC_HINT(nonnull); int fr_control_message_send(fr_control_t *c, void *data, size_t data_size) CC_HINT(nonnull); -ssize_t fr_control_message_receive(fr_atomic_queue_t *aq, struct kevent const *kev, void *data, size_t data_size) CC_HINT(nonnull); +int fr_control_message_service_kevent(fr_atomic_queue_t *aq, struct kevent const *kev) CC_HINT(nonnull); int fr_control_message_push(fr_control_t *c, void *data, size_t data_size) CC_HINT(nonnull); ssize_t fr_control_message_pop(fr_atomic_queue_t *aq, void *data, size_t data_size) CC_HINT(nonnull);