]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
add atomic queue to fr_channel_create()
authorAlan T. DeKok <aland@freeradius.org>
Wed, 7 Dec 2016 21:15:34 +0000 (16:15 -0500)
committerAlan T. DeKok <aland@freeradius.org>
Wed, 7 Dec 2016 22:32:28 +0000 (17:32 -0500)
in preparation for using atomic queues to send data.

src/tests/util/channel_test.c
src/util/channel.c
src/util/channel.h
src/util/worker.c

index d51464f597ad004c2fb1cb1c76236e75ee9acc80..9900986dde337c26eff1d3162abfb6d2f033c741 100644 (file)
@@ -81,7 +81,7 @@ static void *channel_master(void *arg)
        /*
         *      Signal the worker that the channel is open
         */
-       rcode = fr_channel_signal_open(kq_worker, channel);
+       rcode = fr_channel_signal_open(channel);
        if (rcode < 0) {
                fprintf(stderr, "Failed signaling open: %s\n", strerror(errno));
                exit(1);
@@ -189,7 +189,7 @@ check_close:
                        fr_channel_event_t ce;
                        fr_channel_t *new_channel;
 
-                       ce = fr_channel_service_kevent(kq_worker, &events[i], now, &new_channel);
+                       ce = fr_channel_service_kevent(aq_master, &events[i], now, &new_channel);
                        MPRINT1("\tWorker Got channel event %d\n", ce);
 
                        if (ce == FR_CHANNEL_DATA_READY_RECEIVER) {
@@ -216,6 +216,8 @@ check_close:
                                continue;
                        }
 
+                       if (ce == FR_CHANNEL_NOOP) continue;
+
                        fprintf(stderr, "Master got unexpected CE %d\n", ce);
 
                        /*
@@ -299,12 +301,19 @@ static void *channel_worker(void *arg)
                for (i = 0; i < rcode; i++) {
                        fr_channel_event_t ce;
 
-                       ce = fr_channel_service_kevent(kq_worker, &events[i], now, &new_channel);
+                       /*
+                        *      @todo drain the control plane on one kevent.
+                        */
+                       ce = fr_channel_service_kevent(aq_worker, &events[i], now, &new_channel);
                        MPRINT1("\tWorker Got channel event %d\n", ce);
 
                        if (ce == FR_CHANNEL_OPEN) {
                                MPRINT1("\tWorker received a new channel\n");
                                rad_assert(new_channel == channel);
+                               if (fr_channel_worker_receive_open(ctx, new_channel) < 0) {
+                                       fprintf(stderr, "Failed calling receive open.\n");
+                                       exit(1);
+                               }
                                continue;
                        }
 
@@ -371,6 +380,8 @@ static void *channel_worker(void *arg)
                                continue;
                        }
 
+                       if (ce == FR_CHANNEL_NOOP) continue;
+
                        fprintf(stderr, "Worker got unexpected CE %d\n", ce);
 
                        /*
index 4538d16c23f84aab2199b1e21164179b3c09f748..e24c67f7cb01bdfbc1e6c37d1dc42c41dd09cdb5 100644 (file)
@@ -531,7 +531,7 @@ int fr_channel_worker_sleeping(fr_channel_t *ch)
  *  event.  Note that the caller does NOT pass the channel into this
  *  function.  Instead, the channel is taken from the kevent.
  *
- * @param[in] kq the kqueue on which the event was received
+ * @param[in] aq the atomic queue on which we receive control-plane messages
  * @param[in] kev the event of type EVFILT_USER
  * @param[in] when the current time
  * @param[out] p_channel the channel which should be serviced.
@@ -542,7 +542,7 @@ int fr_channel_worker_sleeping(fr_channel_t *ch)
  *     - FR_CHANNEL_OPEN when a channel has been opened and sent to us
  *     - FR_CHANNEL_CLOSE when a channel should be closed
  */
-fr_channel_event_t fr_channel_service_kevent(int kq, struct kevent const *kev, fr_time_t when, fr_channel_t **p_channel)
+fr_channel_event_t fr_channel_service_kevent(UNUSED fr_atomic_queue_t *aq, struct kevent const *kev, fr_time_t when, fr_channel_t **p_channel)
 {
        int rcode;
        uint64_t ack;
@@ -579,8 +579,6 @@ fr_channel_event_t fr_channel_service_kevent(int kq, struct kevent const *kev, f
                 *      worker.  Return the channel to the worker.
                 */
        case FR_CHANNEL_SIGNAL_OPEN:
-               rad_assert(kq == ch->end[TO_WORKER].kq);
-
                *p_channel = ch;
                return FR_CHANNEL_OPEN;
 
@@ -597,18 +595,6 @@ fr_channel_event_t fr_channel_service_kevent(int kq, struct kevent const *kev, f
 
        rad_assert(kev->ident == FR_CHANNEL_SIGNAL_WORKER_SLEEPING);
 
-       /*
-        *      "worker sleeping" signals are only allowed from the
-        *      worker to the master thread.
-        */
-       rad_assert(kq == ch->end[FROM_WORKER].kq);
-
-       /*
-        *      Run-time sanity check.
-        */
-       end = &ch->end[FROM_WORKER];
-       if (end->kq != kq) return FR_CHANNEL_ERROR;
-
        end = &ch->end[TO_WORKER];
 
        /*
@@ -690,23 +676,20 @@ int fr_channel_ack_worker_close(fr_channel_t *ch)
 
 /** Send a channel to a KQ
  *
- * @param[in] kq the kqueue to send the channel
  * @param[in] ch the channel
  * @return
  *     - <0 on error
  *     - 0 on success
  */
-int fr_channel_signal_open(int kq, fr_channel_t *ch)
+int fr_channel_signal_open(fr_channel_t *ch)
 {
        fr_channel_control_t cc;
 
-       rad_assert(kq == ch->end[TO_WORKER].kq);
-
        cc.signal = FR_CHANNEL_SIGNAL_OPEN;
        cc.ack = 0;
        cc.ch = ch;
 
-       return fr_channel_kevent_signal(kq, &cc);
+       return fr_channel_kevent_signal(ch->end[TO_WORKER].kq, &cc);
 }
 
 void fr_channel_debug(fr_channel_t *ch, FILE *fp)
@@ -721,3 +704,31 @@ void fr_channel_debug(fr_channel_t *ch, FILE *fp)
        fprintf(fp, "\tsequence = %zd\n", ch->end[FROM_WORKER].sequence);
        fprintf(fp, "\tack = %zd\n", ch->end[FROM_WORKER].ack);
 }
+
+/** Receive an "open channel" signal.
+ *
+ *  Called only by the worker.
+ *
+ * @param[in] ctx the talloc context for worker messages
+ * @param[in] ch the channel
+ * @return
+ *     - <0 on error
+ *     - 0 on success
+ */
+int fr_channel_worker_receive_open(UNUSED TALLOC_CTX *ctx, UNUSED fr_channel_t *ch)
+{
+#if 0
+#ifndef NDEBUG
+       talloc_get_type_abort(ch, fr_channel_t);
+#endif
+
+       ch->end[FROM_WORKER].control = fr_control_create(ctx,
+                                                        ch->end[FROM_WORKER].kq,
+                                                        ch->end[FROM_WORKER].aq_control);
+       if (!ch->end[FROM_WORKER].control) {
+               return -1;
+       }
+#endif
+
+       return 0;
+}
index 2909a53cee18330a5296ea74e1941447d9d65fc2..6e58c06caafd2b31bf0da8ef235188323fe852b0 100644 (file)
@@ -113,14 +113,17 @@ 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);
 
-fr_channel_event_t fr_channel_service_kevent(int kq, struct kevent const *kev, fr_time_t when, fr_channel_t **p_channel) CC_HINT(nonnull);
+fr_channel_event_t fr_channel_service_kevent(fr_atomic_queue_t *aq, struct kevent const *kev, fr_time_t when, fr_channel_t **p_channel) CC_HINT(nonnull);
 
 bool fr_channel_active(fr_channel_t *ch) CC_HINT(nonnull);
 
-int fr_channel_signal_open(int kq, fr_channel_t *ch) CC_HINT(nonnull);
+int fr_channel_signal_open(fr_channel_t *ch) CC_HINT(nonnull);
+int fr_channel_worker_receive_open(TALLOC_CTX *ctx, fr_channel_t *ch) CC_HINT(nonnull);
+
 int fr_channel_signal_worker_close(fr_channel_t *ch) CC_HINT(nonnull);
 int fr_channel_ack_worker_close(fr_channel_t *ch) CC_HINT(nonnull);
 
+
 void fr_channel_debug(fr_channel_t *ch, FILE *fp);
 
 #ifdef __cplusplus
index 2427c5c65eb02deee1d9e563478bdbbaa2f169c6..2df0e169a8a6cb9955a3f9e25a7fced1b6fa3312 100644 (file)
@@ -34,6 +34,8 @@ RCSID("$Id$")
 struct fr_worker_t {
        int                     kq;             //!< my kq
 
+       fr_atomic_queue_t       *aq_control;    //!< atomic queue for control messages sent to me
+
        fr_message_set_t        *ms;            //!< replies are allocated from here.
 
        fr_event_list_t         *el;            //!< our event list
@@ -75,7 +77,7 @@ struct fr_worker_t {
 /** Handle EVFILT_USER events
  *
  */
-static void fr_worker_evfilt_user(int kq, struct kevent const *kev, void *ctx)
+static void fr_worker_evfilt_user(UNUSED int kq, struct kevent const *kev, void *ctx)
 {
        fr_channel_event_t what;
        fr_channel_t *ch;
@@ -85,7 +87,7 @@ static void fr_worker_evfilt_user(int kq, struct kevent const *kev, void *ctx)
 
        rad_assert(kev->filter == EVFILT_USER);
 
-       what = fr_channel_service_kevent(kq, kev, when, &ch);
+       what = fr_channel_service_kevent(worker->aq_control, kev, when, &ch);
 
        switch (what) {
                /*
@@ -605,6 +607,12 @@ fr_worker_t *fr_worker_create(TALLOC_CTX *ctx, uint32_t num_transports, fr_trans
        worker->kq = fr_event_list_kq(worker->el);
        rad_assert(worker->kq >= 0);
 
+       worker->aq_control = fr_atomic_queue_create(worker, 128);
+       if (!worker->aq_control) {
+               talloc_free(worker);
+               return NULL;
+       }
+
        if (fr_event_user_insert(worker->el, fr_worker_evfilt_user, worker) < 0) {
                talloc_free(worker);
                return NULL;