From: Alan T. DeKok Date: Wed, 7 Dec 2016 19:52:43 +0000 (-0500) Subject: Added atomic queues to fr_channel_create() X-Git-Url: http://git.ipfire.org/cgi-bin/gitweb.cgi?a=commitdiff_plain;h=418eef7b597ff908ae55fc0aae9e68ab21a4e301;p=thirdparty%2Ffreeradius-server.git Added atomic queues to fr_channel_create() --- diff --git a/src/tests/util/channel_test.c b/src/tests/util/channel_test.c index afa08e9b1b2..0254beef43f 100644 --- a/src/tests/util/channel_test.c +++ b/src/tests/util/channel_test.c @@ -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); diff --git a/src/util/channel.c b/src/util/channel.c index 0fceda90c20..9c689808ff5 100644 --- a/src/util/channel.c +++ b/src/util/channel.c @@ -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. diff --git a/src/util/channel.h b/src/util/channel.h index 1e08fea480b..b0e0ef9efc3 100644 --- a/src/util/channel.h +++ b/src/util/channel.h @@ -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);