]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
Added atomic queues to fr_channel_create()
authorAlan T. DeKok <aland@freeradius.org>
Wed, 7 Dec 2016 19:52:43 +0000 (14:52 -0500)
committerAlan T. DeKok <aland@freeradius.org>
Wed, 7 Dec 2016 22:32:27 +0000 (17:32 -0500)
src/tests/util/channel_test.c
src/util/channel.c
src/util/channel.h

index afa08e9b1b270e9f6232718b1b81005ff4ebde48..0254beef43f11ff8958924d6e89203bfa195b917 100644 (file)
@@ -43,6 +43,7 @@ RCSID("$Id$")
 
 static int             debug_lvl = 0;
 static int             kq_master, kq_worker;
+static fr_atomic_queue_t *aq_master, *aq_worker;
 static int             max_messages = 10;
 static int             max_outstanding = 1;
 static bool            touch_memory = false;
@@ -445,7 +446,15 @@ int main(int argc, char *argv[])
        kq_worker = kqueue();
        rad_assert(kq_worker >= 0);
 
-       channel = fr_channel_create(autofree, kq_master, kq_worker);
+       aq_master = fr_atomic_queue_create(autofree, 128);
+       rad_assert(aq_master != NULL);
+
+       aq_worker = fr_atomic_queue_create(autofree, 128);
+       rad_assert(aq_worker != NULL);
+
+       channel = fr_channel_create(autofree,
+                                   kq_master, aq_master,
+                                   kq_worker, aq_worker);
        if (!channel) {
                fprintf(stderr, "channel_test: Failed to create channel\n");
                exit(1);
index 0fceda90c205ec0fc5e6a1853791d962520329cc..9c689808ff51aec50a160ad94fb9a2579cca809d 100644 (file)
@@ -82,6 +82,8 @@ 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
+
        int                     num_outstanding; //!< number of outstanding requests with no reply
 
        size_t                  num_signals;
@@ -154,7 +156,8 @@ static int fr_channel_kevent_signal(int kq, fr_channel_control_t *cc)
  *     - NULL on error
  *     - channel on success
  */
-fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, int kq_worker)
+fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_t *aq_master,
+                               int kq_worker, fr_atomic_queue_t *aq_worker)
 {
        fr_time_t when;
        fr_channel_t *ch;
@@ -177,7 +180,9 @@ fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, int kq_worker)
        }
 
        ch->end[TO_WORKER].kq = kq_worker;
+       ch->end[TO_WORKER].aq_control = aq_worker;
        ch->end[FROM_WORKER].kq = kq_master;
+       ch->end[FROM_WORKER].aq_control = aq_master;
 
        /*
         *      Initialize all of the timers to now.
index 1e08fea480bc1a1ea08dbeca425113d385c8bbef..b0e0ef9efc32842e79aff3f57e4d91778ad2107f 100644 (file)
@@ -102,7 +102,8 @@ typedef enum fr_channel_event_t {
        FR_CHANNEL_CLOSE,
 } fr_channel_event_t;
 
-fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, int kq_worker);
+fr_channel_t *fr_channel_create(TALLOC_CTX *ctx, int kq_master, fr_atomic_queue_t *aq_master,
+                               int kq_worker, fr_atomic_queue_t *aq_worker);
 
 int fr_channel_send_request(fr_channel_t *ch, fr_channel_data_t *cm, fr_channel_data_t **p_reply) CC_HINT(nonnull);
 fr_channel_data_t *fr_channel_recv_request(fr_channel_t *ch) CC_HINT(nonnull);