]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
make fr_channel_service_aq() public
authorAlan T. DeKok <aland@freeradius.org>
Sun, 11 Dec 2016 14:30:06 +0000 (09:30 -0500)
committerAlan T. DeKok <aland@freeradius.org>
Sun, 11 Dec 2016 14:38:00 +0000 (09:38 -0500)
src/tests/util/channel_test.c
src/util/channel.c
src/util/channel.h
src/util/worker.c

index b8ad7a6e4e453e84d248ba1841668e86dc0bf2aa..1149a834677de15183bbfaec5978d93eb2766da9 100644 (file)
@@ -68,6 +68,8 @@ static void *channel_master(void *arg)
        fr_message_set_t *ms;
        TALLOC_CTX *ctx;
        fr_channel_t *channel = arg;
+       fr_channel_t *new_channel;
+       fr_channel_event_t ce;
        struct kevent events[MAX_KEVENTS];
 
        ctx = talloc_init("channel_master");
@@ -172,6 +174,7 @@ static void *channel_master(void *arg)
 check_close:
 
                if (!signaled_close && (num_messages >= max_messages) && (num_outstanding == 0)) {
+                       MPRINT1("Master signaling worker to exit.\n");
                        rcode = fr_channel_signal_worker_close(channel);
                        if (rcode < 0) {
                                fprintf(stderr, "Failed signaling close: %s\n", strerror(errno));
@@ -196,20 +199,22 @@ check_close:
 
                if (num_events == 0) continue;
 
-               now = fr_time();
-
                /*
                 *      Service the events.
                 */
                for (i = 0; i < num_events; i++) {
-                       fr_channel_event_t ce;
-                       fr_channel_t *new_channel;
+                       (void) fr_channel_service_kevent(aq_master, &events[i]);
+               }
 
-                       ce = fr_channel_service_kevent(aq_master, &events[i], now, &new_channel);
-                       MPRINT1("Master got channel event %d\n", ce);
+               now = fr_time();
 
-                       if (ce == FR_CHANNEL_DATA_READY_RECEIVER) {
+               while ((ce = fr_channel_service_aq(aq_master, now, &new_channel)) != FR_CHANNEL_EMPTY) {
+                       MPRINT1("Master got channel event %d\n", ce);
+                       
+                       switch (ce) {
+                       case FR_CHANNEL_DATA_READY_RECEIVER:
                                MPRINT1("Master got data ready signal\n");
+                               rad_assert(new_channel == channel);
 
                                reply = fr_channel_recv_reply(channel);
                                if (!reply) {
@@ -224,33 +229,29 @@ check_close:
                                                num_replies, num_outstanding, num_messages, max_messages);
                                        fr_message_done(&reply->m);
                                } while ((reply = fr_channel_recv_reply(channel)) != NULL);
+                               break;
 
-                               continue;
-                       }
-
-                       if (ce == FR_CHANNEL_CLOSE) {
+                       case FR_CHANNEL_CLOSE:
                                MPRINT1("Master received close signal\n");
+                               rad_assert(new_channel == channel);
                                rad_assert(signaled_close == true);
                                running = false;
-                               continue;
-                       }
-
-                       if (ce == FR_CHANNEL_NOOP) continue;
-
-                       fprintf(stderr, "Master got unexpected CE %d\n", ce);
+                               break;
 
-                       /*
-                        *      Not written yet!
-                        */
-                       rad_assert(0 == 1);
-               }
-       }
+                       case FR_CHANNEL_NOOP:
+                               break;
 
-       MPRINT1("Master signaling worker to exit.\n");
-
-
-       // wait for final set of messages
+                       default:
+                               fprintf(stderr, "Master got unexpected CE %d\n", ce);
 
+                               /*
+                                *      Not written yet!
+                                */
+                               rad_assert(0 == 1);
+                               break;
+                       } /* switch over signal returned */
+               } /* drain the control plane */
+       } /* loop until told to exit */
 
        MPRINT1("Master exiting.\n");
 
@@ -282,6 +283,7 @@ static void *channel_worker(void *arg)
        fr_message_set_t *ms;
        TALLOC_CTX *ctx;
        fr_channel_t *channel = arg;
+       fr_channel_event_t ce;
        struct kevent events[MAX_KEVENTS];
 
        ctx = talloc_init("channel_worker");
@@ -314,18 +316,20 @@ 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]);
+               }
+
                now = fr_time();
 
-               for (i = 0; i < num_events; i++) {
-                       fr_channel_event_t ce;
+               while ((ce = fr_channel_service_aq(aq_worker, now, &new_channel)) != FR_CHANNEL_EMPTY) {
+                       fr_channel_data_t *cd, *reply;
 
-                       /*
-                        *      @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) {
+                       switch (ce) {
+
+                       case FR_CHANNEL_OPEN:
                                MPRINT1("\tWorker received a new channel\n");
                                rad_assert(new_channel == channel);
 
@@ -333,14 +337,11 @@ static void *channel_worker(void *arg)
                                        fprintf(stderr, "Failed calling receive open.\n");
                                        exit(1);
                                }
+                               break;
 
-                               continue;
-                       }
-
-                       if (ce == FR_CHANNEL_CLOSE) {
-                               fr_channel_data_t *cd;
-
+                       case FR_CHANNEL_CLOSE:
                                MPRINT1("\tWorker requested to close the channel.\n");
+                               rad_assert(new_channel == channel);
                                running = false;
 
                                /*
@@ -351,20 +352,18 @@ static void *channel_worker(void *arg)
                                        MPRINT1("\tWorker got message %d\n", worker_messages);
                                        fr_message_done(&cd->m);
                                }
-
+                               
                                (void) fr_channel_ack_worker_close(channel);
-                               continue;
-                       }
-
-                       if (ce == FR_CHANNEL_DATA_READY_WORKER) {
-                               fr_channel_data_t *cd, *reply;
+                               break;
 
+                       case FR_CHANNEL_DATA_READY_WORKER:
                                MPRINT1("\tWorker got data ready signal\n");
+                               rad_assert(new_channel == channel);
 
                                cd = fr_channel_recv_request(channel);
                                if (!cd) {
                                        MPRINT1("\tWorker SIGNAL WITH NO DATA!\n");
-                                       continue;
+                                       break;
                                }
 
                                while (cd) {
@@ -400,18 +399,22 @@ static void *channel_worker(void *arg)
                                        }
                                        rad_assert(rcode == 0);
                                }
-                               continue;
-                       }
+                               break;
 
-                       if (ce == FR_CHANNEL_NOOP) continue;
+                       case FR_CHANNEL_NOOP:
+                               rad_assert(new_channel == channel);
+                               break;
 
-                       fprintf(stderr, "Worker got unexpected CE %d\n", ce);
+                       default:
+                               fprintf(stderr, "Worker got unexpected CE %d\n", ce);
 
-                       /*
-                        *      Not written yet!
-                        */
-                       rad_assert(0 == 1);
-               }
+                               /*
+                                *      Not written yet!
+                                */
+                               rad_assert(0 == 1);
+                               break;
+                       } /* switch over signals */
+               } /* drain the control plane */
        }
 
        MPRINT1("\tWorker exiting.\n");
index c90f897212a964c74bf2ad40a5d69128eacb5759..dc277239258aafcdfd4778124ad4e40b53beb320 100644 (file)
@@ -581,7 +581,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
  */
-static fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t when, fr_channel_t **p_channel)
+fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t when, fr_channel_t **p_channel)
 {
        int rcode;
        uint64_t ack;
@@ -666,16 +666,11 @@ static fr_channel_event_t fr_channel_service_aq(fr_atomic_queue_t *aq, fr_time_t
  *
  * @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.
  * @return
- *     - FR_CHANNEL_ERROR on error
- *     - FR_CHANNEL_NOOP, on do nothing
- *     - FR_CHANNEL_DATA_READY on data ready
- *     - FR_CHANNEL_OPEN when a channel has been opened and sent to us
- *     - FR_CHANNEL_CLOSE when a channel should be closed
+ *     - <0 on error
+ *     - 0 on success
  */
-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)
+int fr_channel_service_kevent(fr_atomic_queue_t *aq, struct kevent const *kev)
 {
        fr_channel_t *ch;
        fr_channel_control_t *cc;
@@ -700,11 +695,10 @@ fr_channel_event_t fr_channel_service_kevent(fr_atomic_queue_t *aq, struct keven
        }
 
        if (!fr_atomic_queue_push(aq, cc)) {
-               *p_channel = NULL;
-               return FR_CHANNEL_ERROR;
+               return -1;
        }
 
-       return fr_channel_service_aq(aq, when, p_channel);
+       return 0;
 }
 
 
index 4e36e76a8cbd54750c62040fdbc4278ff5ad7b20..bbfa8be6a00d1cd9d3017e55fdf07bb6385f66de 100644 (file)
@@ -115,7 +115,8 @@ 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(fr_atomic_queue_t *aq, struct kevent const *kev, fr_time_t when, fr_channel_t **p_channel) CC_HINT(nonnull);
+int fr_channel_service_kevent(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 ced3ca0921e52a9c694662505739caf3ba2781c6..7c20ccdf52943fe483c481c38e2583eb8a5dcb69 100644 (file)
@@ -76,9 +76,11 @@ struct fr_worker_t {
 
 /** Handle EVFILT_USER events
  *
+ * @todo fix for fr_channel_service_aq() also
  */
-static void fr_worker_evfilt_user(UNUSED int kq, struct kevent const *kev, void *ctx)
+static void fr_worker_evfilt_user(UNUSED int kq, UNUSED struct kevent const *kev, UNUSED void *ctx)
 {
+#if 0
        fr_channel_event_t what;
        fr_channel_t *ch;
        fr_channel_data_t *cd;
@@ -132,6 +134,7 @@ static void fr_worker_evfilt_user(UNUSED int kq, struct kevent const *kev, void
        case FR_CHANNEL_ERROR:
                return;
        }
+#endif
 }
 
 /** Decode a request from either the localized queue, or the to_decode queue