]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
move to control-plane signaling for channels
authorAlan T. DeKok <aland@freeradius.org>
Mon, 12 Dec 2016 19:30:14 +0000 (14:30 -0500)
committerAlan T. DeKok <aland@freeradius.org>
Mon, 12 Dec 2016 19:30:14 +0000 (14:30 -0500)
src/tests/util/channel_test.c
src/util/channel.c
src/util/channel.h
src/util/control.c
src/util/control.h

index 05ac568c6de1e3730c821589f5c4959186316594..decc59be2c52970c5a4c3ba5da18d0717761e578 100644 (file)
@@ -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!
index 26ed6b318dc19484d05ae5c5423b0bac872c58e7..27f352870ceda4b7894837800cc163e09425f9d0 100644 (file)
@@ -25,6 +25,7 @@
 RCSID("$Id$")
 
 #include <freeradius-devel/util/channel.h>
+#include <freeradius-devel/util/control.h>
 #include <freeradius-devel/rad_assert.h>
 
 /*
@@ -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;
 }
index bbfa8be6a00d1cd9d3017e55fdf07bb6385f66de..5f926e150e27f928f76ad0b38aa729ec4eb4cb73 100644 (file)
@@ -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);
index 81caad426f9838fab25aa6622b74f77727c0ad84..395975fdb59b046645bf5a95fc3d9bbd2fb12540 100644 (file)
@@ -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;
 }
index 47897a36cea045144a26b4228db3895cdaa4f0c1..341f7c755cd797b7e7295d2c3e14c5554d760abe 100644 (file)
@@ -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);