]> git.ipfire.org Git - thirdparty/freeradius-server.git/commitdiff
renamed receiver to network.
authorAlan T. DeKok <aland@freeradius.org>
Mon, 24 Apr 2017 15:02:09 +0000 (11:02 -0400)
committerAlan T. DeKok <aland@freeradius.org>
Mon, 24 Apr 2017 15:02:09 +0000 (11:02 -0400)
It's not perfect, but it's clearer and makes more sense

src/lib/io/all.mk
src/lib/io/channel.c
src/lib/io/channel.h
src/lib/io/network.c [moved from src/lib/io/receiver.c with 63% similarity]
src/lib/io/network.h [moved from src/lib/io/receiver.h with 62% similarity]
src/lib/io/schedule.c
src/lib/io/worker.c
src/tests/util/channel_test.c
src/tests/util/radius1_test.c
src/tests/util/worker_test.c

index cfb22b5786e80dee55f08c1b9650870fc7c21ea4..4d986c6b5e4726afa351a1727c2f4f32268530b8 100644 (file)
@@ -1,7 +1,7 @@
 TARGET := libfreeradius-io.a
 
 SOURCES        :=      ring_buffer.c message.c atomic_queue.c queue.c time.c channel.c track.c worker.c \
-               schedule.c receiver.c control.c
+               schedule.c network.c control.c
 
 TGT_PREREQS    := libfreeradius-util.la
 TGT_LDLIBS     := $(LIBS)
index c25ff3545c80dfdaa28734060acfebf765e41de8..0c57c1e1d533058f3c8abcadf20851718d9389da 100644 (file)
@@ -65,7 +65,7 @@ RCSID("$Id$")
 typedef enum fr_channel_signal_t {
        FR_CHANNEL_SIGNAL_ERROR                 = FR_CHANNEL_ERROR,
        FR_CHANNEL_SIGNAL_DATA_TO_WORKER        = FR_CHANNEL_DATA_READY_WORKER,
-       FR_CHANNEL_SIGNAL_DATA_FROM_WORKER      = FR_CHANNEL_DATA_READY_RECEIVER,
+       FR_CHANNEL_SIGNAL_DATA_FROM_WORKER      = FR_CHANNEL_DATA_READY_NETWORK,
        FR_CHANNEL_SIGNAL_OPEN                  = FR_CHANNEL_OPEN,
        FR_CHANNEL_SIGNAL_CLOSE                 = FR_CHANNEL_CLOSE,
 
@@ -620,7 +620,7 @@ fr_channel_event_t fr_channel_service_message(fr_time_t when, fr_channel_t **p_c
                 */
        case FR_CHANNEL_SIGNAL_DATA_DONE_WORKER:
                MPRINT("channel got data_done_worker\n");
-               ce = FR_CHANNEL_DATA_READY_RECEIVER;
+               ce = FR_CHANNEL_DATA_READY_NETWORK;
                break;
 
        case FR_CHANNEL_SIGNAL_WORKER_SLEEPING:
index 880ed238f8c9bb28a4f03a4f7a147cecd9ffe566..7585140e0d980edf5b2d2502bb048c1b0730f659 100644 (file)
@@ -47,7 +47,7 @@ typedef struct fr_channel_t fr_channel_t;
 typedef enum fr_channel_event_t {
        FR_CHANNEL_ERROR = 0,
        FR_CHANNEL_DATA_READY_WORKER,
-       FR_CHANNEL_DATA_READY_RECEIVER,
+       FR_CHANNEL_DATA_READY_NETWORK,
        FR_CHANNEL_OPEN,
        FR_CHANNEL_CLOSE,
 
similarity index 63%
rename from src/lib/io/receiver.c
rename to src/lib/io/network.c
index 84276be373046356c6ca77f095ba0e41ed39c259..7dd761dfd87295b02cf3a90b27bbf4ae95478a40 100644 (file)
@@ -18,7 +18,7 @@
  * $Id$
  *
  * @brief Receiver of socket data, which sends messages to the workers.
- * @file io/receiver.c
+ * @file io/network.c
  *
  * @copyright 2016 Alan DeKok <aland@freeradius.org>
  */
@@ -31,20 +31,20 @@ RCSID("$Id$")
 #include <freeradius-devel/io/channel.h>
 #include <freeradius-devel/io/control.h>
 #include <freeradius-devel/io/worker.h>
-#include <freeradius-devel/io/receiver.h>
+#include <freeradius-devel/io/network.h>
 
 #include <freeradius-devel/rad_assert.h>
 
-typedef struct fr_receiver_worker_t {
+typedef struct fr_network_worker_t {
        int                     heap_id;                //!< workers are in a heap
        fr_time_t               cpu_time;               //!< how much CPU time this worker has spent
        fr_time_t               predicted;              //!< predicted processing time for one packet
 
        fr_channel_t            *channel;               //!< channel to the worker
        fr_worker_t             *worker;                //!< worker pointer
-} fr_receiver_worker_t;
+} fr_network_worker_t;
 
-typedef struct fr_receiver_socket_t {
+typedef struct fr_network_socket_t {
        int                     heap_id;                //!< for the heap
 
        int                     fd;                     //!< the file descriptor
@@ -53,10 +53,10 @@ typedef struct fr_receiver_socket_t {
 
        fr_message_set_t        *ms;                    //!< message buffers for this socket.
        fr_channel_data_t       *cd;                    //!< cached in case of allocation & read error
-} fr_receiver_socket_t;
+} fr_network_socket_t;
 
 
-struct fr_receiver_t {
+struct fr_network_t {
        int                     kq;                     //!< our KQ
 
        fr_log_t                *log;                   //!< log destination
@@ -85,8 +85,8 @@ struct fr_receiver_t {
 
 static int socket_cmp(void const *one, void const *two)
 {
-       fr_receiver_socket_t const *a = one;
-       fr_receiver_socket_t const *b = two;
+       fr_network_socket_t const *a = one;
+       fr_network_socket_t const *b = two;
 
        if (a->fd < b->fd) return -1;
        if (a->fd > b->fd) return +1;
@@ -96,8 +96,8 @@ static int socket_cmp(void const *one, void const *two)
 
 static int worker_cmp(void const *one, void const *two)
 {
-       fr_receiver_worker_t const *a = one;
-       fr_receiver_worker_t const *b = two;
+       fr_network_worker_t const *a = one;
+       fr_network_worker_t const *b = two;
 
        if (a->cpu_time < b->cpu_time) return -1;
        if (a->cpu_time > b->cpu_time) return +1;
@@ -124,13 +124,13 @@ static int reply_cmp(void const *one, void const *two)
 
 /** Drain the input channel
  *
- * @param[in] rc the receiver
+ * @param[in] nr the network
  * @param[in] ch the channel to drain
  * @param[in] cd the message (if any) to start with
  */
-static void fr_receiver_drain_input(fr_receiver_t *rc, fr_channel_t *ch, fr_channel_data_t *cd)
+static void fr_network_drain_input(fr_network_t *nr, fr_channel_t *ch, fr_channel_data_t *cd)
 {
-       fr_receiver_worker_t *w;
+       fr_network_worker_t *w;
 
        if (!cd) {
                cd = fr_channel_recv_reply(ch);
@@ -140,8 +140,8 @@ static void fr_receiver_drain_input(fr_receiver_t *rc, fr_channel_t *ch, fr_chan
        }
 
        do {
-               rc->num_replies++;
-               fr_log(rc->log, L_DBG, "received reply %zd", rc->num_replies);
+               nr->num_replies++;
+               fr_log(nr->log, L_DBG, "received reply %zd", nr->num_replies);
 
                cd->channel.ch = ch;
 
@@ -156,7 +156,7 @@ static void fr_receiver_drain_input(fr_receiver_t *rc, fr_channel_t *ch, fr_chan
                        w->predicted = RTT(w->predicted, cd->reply.processing_time);
                }
 
-               (void) fr_heap_insert(rc->replies, cd);
+               (void) fr_heap_insert(nr->replies, cd);
        } while ((cd = fr_channel_recv_reply(ch)) != NULL);
 }
 
@@ -166,14 +166,14 @@ static void fr_receiver_drain_input(fr_receiver_t *rc, fr_channel_t *ch, fr_chan
  *  work, and tell the event code to return to the main loop if
  *  there's work to do.
  *
- * @param[in] ctx the receiver
+ * @param[in] ctx the network
  * @param[in] wake the time when the event loop will wake up.
  */
-static int fr_receiver_idle(void *ctx, struct timeval *wake)
+static int fr_network_idle(void *ctx, struct timeval *wake)
 {
-       fr_receiver_t *rc = talloc_get_type_abort(ctx, fr_receiver_t);
+       fr_network_t *nr = talloc_get_type_abort(ctx, fr_network_t);
 
-       rad_cond_assert(rc->el != NULL); /* temporary until we actually use rc here */
+       rad_cond_assert(nr->el != NULL); /* temporary until we actually use nr here */
 
        if (!wake) {
                // ready to process requests
@@ -193,51 +193,51 @@ static int fr_receiver_idle(void *ctx, struct timeval *wake)
 }
 
 
-/** Handle a receiver control message callback for a channel
+/** Handle a network control message callback for a channel
  *
- * @param[in] ctx the receiver
+ * @param[in] ctx the network
  * @param[in] data the message
  * @param[in] data_size size of the data
  * @param[in] now the current time
  */
-static void fr_receiver_channel_callback(void *ctx, void const *data, size_t data_size, fr_time_t now)
+static void fr_network_channel_callback(void *ctx, void const *data, size_t data_size, fr_time_t now)
 {
        fr_channel_event_t ce;
        fr_channel_t *ch;
-       fr_receiver_t *rc = ctx;
+       fr_network_t *nr = ctx;
 
        ce = fr_channel_service_message(now, &ch, data, data_size);
        switch (ce) {
        case FR_CHANNEL_ERROR:
-               fr_log(rc->log, L_DBG_ERR, "aq error");
+               fr_log(nr->log, L_DBG_ERR, "aq error");
                return;
 
        case FR_CHANNEL_EMPTY:
-               fr_log(rc->log, L_DBG, "aq empty");
+               fr_log(nr->log, L_DBG, "aq empty");
                return;
 
        case FR_CHANNEL_NOOP:
-               fr_log(rc->log, L_DBG, "aq noop");
+               fr_log(nr->log, L_DBG, "aq noop");
                break;
 
-       case FR_CHANNEL_DATA_READY_RECEIVER:
+       case FR_CHANNEL_DATA_READY_NETWORK:
                rad_assert(ch != NULL);
-               fr_log(rc->log, L_DBG, "aq data ready");
-               fr_receiver_drain_input(rc, ch, NULL);
+               fr_log(nr->log, L_DBG, "aq data ready");
+               fr_network_drain_input(nr, ch, NULL);
                break;
 
        case FR_CHANNEL_DATA_READY_WORKER:
                rad_assert(0 == 1);
-               fr_log(rc->log, L_DBG_ERR, "aq data ready ? WORKER ?");
+               fr_log(nr->log, L_DBG_ERR, "aq data ready ? WORKER ?");
                break;
 
        case FR_CHANNEL_OPEN:
                rad_assert(0 == 1);
-               fr_log(rc->log, L_DBG, "channel open ?");
+               fr_log(nr->log, L_DBG, "channel open ?");
                break;
 
        case FR_CHANNEL_CLOSE:
-               fr_log(rc->log, L_DBG, "aq channel close");
+               fr_log(nr->log, L_DBG, "aq channel close");
                ///
                break;
        }
@@ -245,26 +245,26 @@ static void fr_receiver_channel_callback(void *ctx, void const *data, size_t dat
 
 /** Send a message on the "best" channel.
  *
- * @param rc the receiver
+ * @param nr the network
  * @param cd the message we've received
  */
-static int fr_receiver_send_request(fr_receiver_t *rc, fr_channel_data_t *cd)
+static int fr_network_send_request(fr_network_t *nr, fr_channel_data_t *cd)
 {
-       fr_receiver_worker_t *worker;
+       fr_network_worker_t *worker;
        fr_channel_data_t *reply;
 
-       (void) talloc_get_type_abort(rc, fr_receiver_t);
+       (void) talloc_get_type_abort(nr, fr_network_t);
 
        /*
         *      Grab the worker with the least total CPU time.
         */
-       worker = fr_heap_pop(rc->workers);
+       worker = fr_heap_pop(nr->workers);
        if (!worker) {
-               fr_log(rc->log, L_DBG, "no workers");
+               fr_log(nr->log, L_DBG, "no workers");
                return 0;
        }
 
-       (void) talloc_get_type_abort(worker, fr_receiver_worker_t);
+       (void) talloc_get_type_abort(worker, fr_network_worker_t);
 
        /*
         *      Send the message to the channel.  If we fail, recurse.
@@ -282,8 +282,8 @@ static int fr_receiver_send_request(fr_receiver_t *rc, fr_channel_data_t *cd)
        if (fr_channel_send_request(worker->channel, cd, &reply) < 0) {
                int rcode;
 
-               fr_log(rc->log, L_DBG, "recursing in send_request");
-               rcode = fr_receiver_send_request(rc, cd);
+               fr_log(nr->log, L_DBG, "recursing in send_request");
+               rcode = fr_network_send_request(nr, cd);
 
                /*
                 *      Mark this channel as still busy, for some
@@ -292,7 +292,7 @@ static int fr_receiver_send_request(fr_receiver_t *rc, fr_channel_data_t *cd)
                 *      to send it another request.
                 */
                worker->cpu_time = cd->m.when + worker->predicted;
-               (void) fr_heap_insert(rc->workers, worker);
+               (void) fr_heap_insert(nr->workers, worker);
 
                return rcode;
        }
@@ -308,13 +308,13 @@ static int fr_receiver_send_request(fr_receiver_t *rc, fr_channel_data_t *cd)
        /*
         *      Insert the worker back into the heap of workers.
         */
-       (void) fr_heap_insert(rc->workers, worker);
+       (void) fr_heap_insert(nr->workers, worker);
 
        /*
         *      If we have a reply, push it onto our local queue, and
         *      poll for more replies.
         */
-       if (reply) fr_receiver_drain_input(rc, worker->channel, reply);
+       if (reply) fr_network_drain_input(nr, worker->channel, reply);
 
        return 1;
 }
@@ -327,23 +327,23 @@ static fr_time_t start_time = 0;
  *
  * @param el the event list
  * @param sockfd the socket which is ready to read
- * @param ctx the receiver socket context.
+ * @param ctx the network socket context.
  */
-static void fr_receiver_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
+static void fr_network_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
 {
-       fr_receiver_socket_t *s = ctx;
-       fr_receiver_t *rc = talloc_parent(s);
+       fr_network_socket_t *s = ctx;
+       fr_network_t *nr = talloc_parent(s);
        ssize_t data_size;
        fr_channel_data_t *cd;
 
        rad_assert(s->fd == sockfd);
 
-       fr_log(rc->log, L_DBG, "receiver read");
+       fr_log(nr->log, L_DBG, "network read");
 
        if (!s->cd) {
                cd = (fr_channel_data_t *) fr_message_reserve(s->ms, s->transport->default_message_size);
                if (!cd) {
-                       fr_log(rc->log, L_ERR, "Failed allocating message size %zd!", s->transport->default_message_size);
+                       fr_log(nr->log, L_ERR, "Failed allocating message size %zd!", s->transport->default_message_size);
 
                        /*
                         *      @todo - handle errors via transport callback
@@ -362,7 +362,7 @@ static void fr_receiver_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
         */
        data_size = s->transport->read(sockfd, s->ctx, cd->m.data, cd->m.rb_size);
        if (data_size == 0) {
-               fr_log(rc->log, L_DBG_ERR, "got no data from transport read");
+               fr_log(nr->log, L_DBG_ERR, "got no data from transport read");
 
                s->cd = cd;
 
@@ -372,7 +372,7 @@ static void fr_receiver_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
        }
 
        if (data_size < 0) {
-               fr_log(rc->log, L_DBG_ERR, "error from transport read");
+               fr_log(nr->log, L_DBG_ERR, "error from transport read");
 
                /*
                 *      @todo - handle errors via transport callback
@@ -381,7 +381,7 @@ static void fr_receiver_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
        }
        s->cd = NULL;
 
-       fr_log(rc->log, L_DBG, "got packet size %zd", data_size);
+       fr_log(nr->log, L_DBG, "got packet size %zd", data_size);
 
        /*
         *      Initialize the rest of the fields of the channel data.
@@ -397,29 +397,29 @@ static void fr_receiver_read(UNUSED fr_event_list_t *el, int sockfd, void *ctx)
 
        (void) fr_message_alloc(s->ms, &cd->m, data_size);
 
-       if (!fr_receiver_send_request(rc, cd)) {
-               fr_log(rc->log, L_ERR, "Failed sending packet to worker");
+       if (!fr_network_send_request(nr, cd)) {
+               fr_log(nr->log, L_ERR, "Failed sending packet to worker");
                fr_message_done(&cd->m);
        }
 }
 
-/** Handle a receiver control message callback for a new socket
+/** Handle a network control message callback for a new socket
  *
- * @param[in] ctx the receiver
+ * @param[in] ctx the network
  * @param[in] data the message
  * @param[in] data_size size of the data
  * @param[in] now the current time
  */
-static void fr_receiver_socket_callback(void *ctx, void const *data, size_t data_size, UNUSED fr_time_t now)
+static void fr_network_socket_callback(void *ctx, void const *data, size_t data_size, UNUSED fr_time_t now)
 {
-       fr_receiver_t *rc = ctx;
-       fr_receiver_socket_t *s;
+       fr_network_t *nr = ctx;
+       fr_network_socket_t *s;
 
        rad_assert(data_size == sizeof(*s));
 
        if (data_size != sizeof(*s)) return;
 
-       s = talloc(rc, fr_receiver_socket_t);
+       s = talloc(nr, fr_network_socket_t);
        rad_assert(s != NULL);
        memcpy(s, data, sizeof(*s));
 
@@ -432,7 +432,7 @@ static void fr_receiver_socket_callback(void *ctx, void const *data, size_t data
                                      sizeof(fr_channel_data_t),
                                      s->transport->default_message_size * MIN_MESSAGES);
        if (!s->ms) {
-               fr_log(rc->log, L_ERR, "Failed creating message buffers for network IO.");
+               fr_log(nr->log, L_ERR, "Failed creating message buffers for network IO.");
 
                /*
                 *      @todo - handle errors via transport callback
@@ -440,46 +440,46 @@ static void fr_receiver_socket_callback(void *ctx, void const *data, size_t data
                _exit(1);
        }
 
-       if (fr_event_fd_insert(rc->el, s->fd, fr_receiver_read, NULL, NULL, s) < 0) {
-               fr_log(rc->log, L_ERR, "Failed adding new socket to event loop: %s", fr_strerror());
+       if (fr_event_fd_insert(nr->el, s->fd, fr_network_read, NULL, NULL, s) < 0) {
+               fr_log(nr->log, L_ERR, "Failed adding new socket to event loop: %s", fr_strerror());
                close(s->fd);
                return;
        }
 
-       (void) fr_heap_insert(rc->sockets, s);
+       (void) fr_heap_insert(nr->sockets, s);
 
-       fr_log(rc->log, L_DBG, "Using new socket with FD %d", s->fd);
+       fr_log(nr->log, L_DBG, "Using new socket with FD %d", s->fd);
 }
 
 
-/** Handle a receiver control message callback for a new worker
+/** Handle a network control message callback for a new worker
  *
- * @param[in] ctx the receiver
+ * @param[in] ctx the network
  * @param[in] data the message
  * @param[in] data_size size of the data
  * @param[in] now the current time
  */
-static void fr_receiver_worker_callback(void *ctx, void const *data, size_t data_size, UNUSED fr_time_t now)
+static void fr_network_worker_callback(void *ctx, void const *data, size_t data_size, UNUSED fr_time_t now)
 {
-       fr_receiver_t *rc = ctx;
+       fr_network_t *nr = ctx;
        fr_worker_t *worker;
-       fr_receiver_worker_t *w;
+       fr_network_worker_t *w;
 
        rad_assert(data_size == sizeof(worker));
 
        memcpy(&worker, data, data_size);
        (void) talloc_get_type_abort(worker, fr_worker_t);
 
-       w = talloc_zero(rc, fr_receiver_worker_t);
+       w = talloc_zero(nr, fr_network_worker_t);
        if (!w) _exit(1);
 
        w->worker = worker;
-       w->channel = fr_worker_channel_create(worker, w, rc->control);
+       w->channel = fr_worker_channel_create(worker, w, nr->control);
        if (!w->channel) _exit(1);
 
        fr_channel_master_ctx_add(w->channel, w);
 
-       (void) fr_heap_insert(rc->workers, w);
+       (void) fr_heap_insert(nr->workers, w);
 }
 
 
@@ -489,14 +489,14 @@ static void fr_receiver_worker_callback(void *ctx, void const *data, size_t data
  * @param[in] kev the kevent to service
  * @param[in] ctx the fr_worker_t
  */
-static void fr_receiver_evfilt_user(UNUSED int kq, struct kevent const *kev, void *ctx)
+static void fr_network_evfilt_user(UNUSED int kq, struct kevent const *kev, void *ctx)
 {
        fr_time_t now;
-       fr_receiver_t *rc = talloc_get_type_abort(ctx, fr_receiver_t);
+       fr_network_t *nr = talloc_get_type_abort(ctx, fr_network_t);
        uint8_t data[256];
 
-       if (!fr_control_message_service_kevent(rc->control, kev)) {
-               fr_log(rc->log, L_DBG, "kevent not for us: ignoring");
+       if (!fr_control_message_service_kevent(nr->control, kev)) {
+               fr_log(nr->log, L_DBG, "kevent not for us: ignoring");
                return;
        }
 
@@ -505,11 +505,11 @@ static void fr_receiver_evfilt_user(UNUSED int kq, struct kevent const *kev, voi
        /*
         *      Service all available control-plane events
         */
-       fr_control_service(rc->control, data, sizeof(data), now);
+       fr_control_service(nr->control, data, sizeof(data), now);
 }
 
 
-/** Create a receiver
+/** Create a network
  *
  * @param[in] ctx the talloc ctx
  * @param[in] logger the destination for all logging messages
@@ -517,135 +517,135 @@ static void fr_receiver_evfilt_user(UNUSED int kq, struct kevent const *kev, voi
  * @param[in] transports the array of transports.
  * @return
  *     - NULL on error
- *     - fr_receiver_t on success
+ *     - fr_network_t on success
  */
-fr_receiver_t *fr_receiver_create(TALLOC_CTX *ctx, fr_log_t *logger, uint32_t num_transports, fr_transport_t **transports)
+fr_network_t *fr_network_create(TALLOC_CTX *ctx, fr_log_t *logger, uint32_t num_transports, fr_transport_t **transports)
 {
-       fr_receiver_t *rc;
+       fr_network_t *nr;
 
        if (!num_transports || !transports) {
                fr_strerror_printf("Must specify a transport");
                return NULL;
        }
 
-       rc = talloc_zero(ctx, fr_receiver_t);
-       if (!rc) {
+       nr = talloc_zero(ctx, fr_network_t);
+       if (!nr) {
        nomem:
                fr_strerror_printf("Failed allocating memory");
                return NULL;
        }
 
-       rc->el = fr_event_list_create(rc, fr_receiver_idle, rc);
-       if (!rc->el) {
+       nr->el = fr_event_list_create(nr, fr_network_idle, nr);
+       if (!nr->el) {
                fr_strerror_printf("Failed creating event list: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       rc->log = logger;
+       nr->log = logger;
 
-       rc->kq = fr_event_list_kq(rc->el);
-       rad_assert(rc->kq >= 0);
+       nr->kq = fr_event_list_kq(nr->el);
+       rad_assert(nr->kq >= 0);
 
-       rc->aq_control = fr_atomic_queue_create(rc, 1024);
-       if (!rc->aq_control) {
-               talloc_free(rc);
+       nr->aq_control = fr_atomic_queue_create(nr, 1024);
+       if (!nr->aq_control) {
+               talloc_free(nr);
                goto nomem;
        }
 
-       rc->control = fr_control_create(rc, rc->kq, rc->aq_control);
-       if (!rc->control) {
+       nr->control = fr_control_create(nr, nr->kq, nr->aq_control);
+       if (!nr->control) {
                fr_strerror_printf("Failed creating control queue: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       rc->rb = fr_ring_buffer_create(rc, FR_CONTROL_MAX_MESSAGES * FR_CONTROL_MAX_SIZE);
-       if (!rc->rb) {
+       nr->rb = fr_ring_buffer_create(nr, FR_CONTROL_MAX_MESSAGES * FR_CONTROL_MAX_SIZE);
+       if (!nr->rb) {
                fr_strerror_printf("Failed creating ring buffer: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       if (fr_control_callback_add(rc->control, FR_CONTROL_ID_CHANNEL, rc, fr_receiver_channel_callback) < 0) {
+       if (fr_control_callback_add(nr->control, FR_CONTROL_ID_CHANNEL, nr, fr_network_channel_callback) < 0) {
                fr_strerror_printf("Failed adding channel callback: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       if (fr_control_callback_add(rc->control, FR_CONTROL_ID_SOCKET, rc, fr_receiver_socket_callback) < 0) {
+       if (fr_control_callback_add(nr->control, FR_CONTROL_ID_SOCKET, nr, fr_network_socket_callback) < 0) {
                fr_strerror_printf("Failed adding socket callback: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       if (fr_control_callback_add(rc->control, FR_CONTROL_ID_WORKER, rc, fr_receiver_worker_callback) < 0) {
+       if (fr_control_callback_add(nr->control, FR_CONTROL_ID_WORKER, nr, fr_network_worker_callback) < 0) {
                fr_strerror_printf("Failed adding worker callback: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
-       if (fr_event_user_insert(rc->el, fr_receiver_evfilt_user, rc) < 0) {
+       if (fr_event_user_insert(nr->el, fr_network_evfilt_user, nr) < 0) {
                fr_strerror_printf("Failed updating event list: %s", fr_strerror());
-               talloc_free(rc);
+               talloc_free(nr);
                return NULL;
        }
 
        /*
         *      Create the various heaps.
         */
-       rc->sockets = fr_heap_create(socket_cmp, offsetof(fr_receiver_socket_t, heap_id));
-       if (!rc->sockets) {
-               talloc_free(rc);
+       nr->sockets = fr_heap_create(socket_cmp, offsetof(fr_network_socket_t, heap_id));
+       if (!nr->sockets) {
+               talloc_free(nr);
                goto nomem;
        }
 
-       rc->replies = fr_heap_create(reply_cmp, offsetof(fr_channel_data_t, channel.heap_id));
-       if (!rc->replies) {
-               talloc_free(rc);
+       nr->replies = fr_heap_create(reply_cmp, offsetof(fr_channel_data_t, channel.heap_id));
+       if (!nr->replies) {
+               talloc_free(nr);
                goto nomem;
        }
 
-       rc->workers = fr_heap_create(worker_cmp, offsetof(fr_channel_data_t, channel.heap_id));
-       if (!rc->workers) {
-               talloc_free(rc);
+       nr->workers = fr_heap_create(worker_cmp, offsetof(fr_channel_data_t, channel.heap_id));
+       if (!nr->workers) {
+               talloc_free(nr);
                goto nomem;
        }
 
-       rc->closing = fr_heap_create(worker_cmp, offsetof(fr_channel_data_t, channel.heap_id));
-       if (!rc->closing) {
-               talloc_free(rc);
+       nr->closing = fr_heap_create(worker_cmp, offsetof(fr_channel_data_t, channel.heap_id));
+       if (!nr->closing) {
+               talloc_free(nr);
                goto nomem;
        }
 
-       rc->num_transports = num_transports;
-       rc->transports = transports;
+       nr->num_transports = num_transports;
+       nr->transports = transports;
 
-       return rc;
+       return nr;
 }
 
 
-/** Destroy a receiver
+/** Destroy a network
  *
- * @param[in] rc the receiver
+ * @param[in] nr the network
  * @return
  *     - <0 on error
  *     - 0 on success
  */
-int fr_receiver_destroy(fr_receiver_t *rc)
+int fr_network_destroy(fr_network_t *nr)
 {
-       fr_receiver_worker_t *worker;
+       fr_network_worker_t *worker;
        fr_channel_data_t *cd;
 
-       (void) talloc_get_type_abort(rc, fr_receiver_t);
+       (void) talloc_get_type_abort(nr, fr_network_t);
 
        /*
         *      Pop all of the workers, and signal them that we're
         *      closing/
         */
-       while ((worker = fr_heap_pop(rc->workers)) != NULL) {
+       while ((worker = fr_heap_pop(nr->workers)) != NULL) {
                fr_channel_signal_worker_close(worker->channel);
-               (void) fr_heap_insert(rc->closing, worker);
+               (void) fr_heap_insert(nr->closing, worker);
        }
 
        /*
@@ -659,54 +659,54 @@ int fr_receiver_destroy(fr_receiver_t *rc)
         *      @todo - call transport "done" for the reply, so that
         *      it knows the replies are done, too.
         */
-       while ((cd = fr_heap_pop(rc->replies)) != NULL) {
+       while ((cd = fr_heap_pop(nr->replies)) != NULL) {
                fr_message_done(&cd->m);
        }
 
-       talloc_free(rc);
+       talloc_free(nr);
 
        return 0;
 }
 
 /** The main network worker function.
  *
- * @param[in] rc the receiver data structure to run.
+ * @param[in] nr the network data structure to run.
  */
-void fr_receiver(fr_receiver_t *rc)
+void fr_network(fr_network_t *nr)
 {
        while (true) {
                bool wait_for_event;
                int num_events;
 //             fr_time_t now;
                fr_channel_data_t *cd;
-               fr_receiver_socket_t *s;
+               fr_network_socket_t *s;
 
                /*
                 *      There are runnable requests.  We still service
                 *      the event loop, but we don't wait for events.
                 */
-               wait_for_event = (fr_heap_num_elements(rc->replies) == 0);
-               fr_log(rc->log, L_DBG, "Waiting for events %d", wait_for_event);
+               wait_for_event = (fr_heap_num_elements(nr->replies) == 0);
+               fr_log(nr->log, L_DBG, "Waiting for events %d", wait_for_event);
 
                /*
                 *      Check the event list.  If there's an error
                 *      (e.g. exit), we stop looping and clean up.
                 */
-               num_events = fr_event_corral(rc->el, wait_for_event);
-               fr_log(rc->log, L_DBG, "Got num_events %d", num_events);
+               num_events = fr_event_corral(nr->el, wait_for_event);
+               fr_log(nr->log, L_DBG, "Got num_events %d", num_events);
                if (num_events < 0) break;
 
                /*
                 *      Service outstanding events.
                 */
                if (num_events > 0) {
-                       fr_log(rc->log, L_DBG, "servicing events");
-                       fr_event_service(rc->el);
+                       fr_log(nr->log, L_DBG, "servicing events");
+                       fr_event_service(nr->el);
                }
 
 //             now = fr_time();
 
-               cd = fr_heap_pop(rc->replies);
+               cd = fr_heap_pop(nr->replies);
                if (!cd) continue;
 
                /*
@@ -716,7 +716,7 @@ void fr_receiver(fr_receiver_t *rc)
 
                s->transport->write(s->fd, s->ctx, cd->m.data, cd->m.data_size);
 
-               fr_log(rc->log, L_DBG, "handling reply to socket %p", cd->io_ctx);
+               fr_log(nr->log, L_DBG, "handling reply to socket %p", cd->io_ctx);
                fr_message_done(&cd->m);
        }
 }
@@ -725,41 +725,41 @@ void fr_receiver(fr_receiver_t *rc)
  *
  *  WARNING: This may be called from another thread!  Care is required.
  *
- * @param[in] rc the receiver data structure to manage
+ * @param[in] nr the network data structure to manage
  */
-void fr_receiver_exit(fr_receiver_t *rc)
+void fr_network_exit(fr_network_t *nr)
 {
-       fr_event_loop_exit(rc->el, 1);
+       fr_event_loop_exit(nr->el, 1);
 }
 
-/** Add a socket to a receiver
+/** Add a socket to a network
  *
- * @param rc the receiver
+ * @param nr the network
  * @param fd the file descriptor for the socket
  * @param ctx the context for the transport
  * @param transport the transport
  */
-int fr_receiver_socket_add(fr_receiver_t *rc, int fd, void *ctx, fr_transport_t *transport)
+int fr_network_socket_add(fr_network_t *nr, int fd, void *ctx, fr_transport_t *transport)
 {
-       fr_receiver_socket_t m;
+       fr_network_socket_t m;
 
        memset(&m, 0, sizeof(m));
        m.fd = fd;
        m.ctx = ctx;
        m.transport = transport;
 
-       return fr_control_message_send(rc->control, rc->rb, FR_CONTROL_ID_SOCKET, &m, sizeof(m));
+       return fr_control_message_send(nr->control, nr->rb, FR_CONTROL_ID_SOCKET, &m, sizeof(m));
 }
 
-/** Add a worker to a receiver
+/** Add a worker to a network
  *
- * @param rc the receiver
+ * @param nr the network
  * @param worker the worker
  */
-int fr_receiver_worker_add(fr_receiver_t *rc, fr_worker_t *worker)
+int fr_network_worker_add(fr_network_t *nr, fr_worker_t *worker)
 {
-       (void) talloc_get_type_abort(rc, fr_receiver_t);
+       (void) talloc_get_type_abort(nr, fr_network_t);
        (void) talloc_get_type_abort(worker, fr_worker_t);
 
-       return fr_control_message_send(rc->control, rc->rb, FR_CONTROL_ID_WORKER, &worker, sizeof(worker));
+       return fr_control_message_send(nr->control, nr->rb, FR_CONTROL_ID_WORKER, &worker, sizeof(worker));
 }
similarity index 62%
rename from src/lib/io/receiver.h
rename to src/lib/io/network.h
index dcc655646f2e5c980d41903431dcecb7a339ad78..002d88ed48b1b217244370c7f67ea7e2edef097a 100644 (file)
 /**
  * $Id$
  *
- * @file io/receiver.h
+ * @file io/network.h
  * @brief Receive packets
  *
  * @copyright 2016 Alan DeKok <aland@freeradius.org>
  */
-RCSIDH(receiver_h, "$Id$")
+RCSIDH(network_h, "$Id$")
 
 #include <freeradius-devel/fr_log.h>
 
@@ -31,18 +31,18 @@ RCSIDH(receiver_h, "$Id$")
 extern "C" {
 #endif
 
-typedef struct fr_receiver_t fr_receiver_t;
+typedef struct fr_network_t fr_network_t;
 
-fr_receiver_t *fr_receiver_create(TALLOC_CTX *ctx, fr_log_t *logger, uint32_t num_transports, fr_transport_t **transports);
-void fr_receiver_exit(fr_receiver_t *rc);
-int fr_receiver_destroy(fr_receiver_t *rc) CC_HINT(nonnull);
-void fr_receiver(fr_receiver_t *rc) CC_HINT(nonnull);
+fr_network_t *fr_network_create(TALLOC_CTX *ctx, fr_log_t *logger, uint32_t num_transports, fr_transport_t **transports);
+void fr_network_exit(fr_network_t *nr);
+int fr_network_destroy(fr_network_t *nr) CC_HINT(nonnull);
+void fr_network(fr_network_t *nr) CC_HINT(nonnull);
 
-int fr_receiver_socket_add(fr_receiver_t *rc, int fd, void *ctx, fr_transport_t *transport) CC_HINT(nonnull);
-int fr_receiver_worker_add(fr_receiver_t *rc, fr_worker_t *worker) CC_HINT(nonnull);
+int fr_network_socket_add(fr_network_t *nr, int fd, void *ctx, fr_transport_t *transport) CC_HINT(nonnull);
+int fr_network_worker_add(fr_network_t *nr, fr_worker_t *worker) CC_HINT(nonnull);
 
 #ifdef __cplusplus
 }
 #endif
 
-#endif /* _FR_RECEIVER_H */
+#endif /* _FR_NETWORK_H */
index 4fb71ceb3e57a29daeaa73b1aa6351c207f4f0e0..fcba83344a9e6a05faf67356f23b441d68f6bddb 100644 (file)
@@ -30,7 +30,7 @@ RCSID("$Id$")
 #include <freeradius-devel/io/schedule.h>
 #include <freeradius-devel/rbtree.h>
 
-#include <freeradius-devel/io/receiver.h>
+#include <freeradius-devel/io/network.h>
 
 #ifdef HAVE_PTHREAD_H
 #include <pthread.h>
@@ -99,17 +99,17 @@ typedef struct fr_schedule_worker_t {
 } fr_schedule_worker_t;
 
 /**
- *     A data structure to track network threads / receivers.
+ *     A data structure to track network threads / networks.
  */
-typedef struct fr_schedule_receiver_t {
-       pthread_t       pthread_id;             //!< the thread of this receiver
+typedef struct fr_schedule_network_t {
+       pthread_t       pthread_id;             //!< the thread of this network
 
        int             id;                     //!< a unique ID
        fr_schedule_t   *sc;                    //!< the scheduler we are running under
 
        fr_schedule_child_status_t status;      //!< status of the worker
-       fr_receiver_t   *rc;                    //!< the receive data structure
-} fr_schedule_receiver_t;
+       fr_network_t    *rc;                    //!< the receive data structure
+} fr_schedule_network_t;
 
 
 /**
@@ -138,7 +138,7 @@ struct fr_schedule_t {
        fr_heap_t       *workers;               //!< heap of workers
        fr_heap_t       *done_workers;          //!< heap of done workers
 
-       fr_schedule_receiver_t *sr;             //!< pointer to the (one) network thread
+       fr_schedule_network_t *sn;              //!< pointer to the (one) network thread
 
        uint32_t        num_transports;         //!< how many transport layers we have
        fr_transport_t  **transports;           //!< array of active transports.
@@ -237,7 +237,7 @@ static void *fr_schedule_worker_thread(void *arg)
        sc->num_workers++;
        PTHREAD_MUTEX_UNLOCK(&sc->mutex);
 
-       (void) fr_receiver_worker_add(sc->sr->rc, sw->worker);
+       (void) fr_network_worker_add(sc->sn->rc, sw->worker);
 
        fr_log(sc->log, L_INFO, "Worker %d running\n", sw->id);
 
@@ -295,33 +295,33 @@ fail:
 }
 
 
-/** Initialize and run the receiver thread.
+/** Initialize and run the network thread.
  *
- * @param[in] arg the fr_schedule_receiver_t
+ * @param[in] arg the fr_schedule_network_t
  * @return NULL
  */
-static void *fr_schedule_receiver_thread(void *arg)
+static void *fr_schedule_network_thread(void *arg)
 {
        TALLOC_CTX *ctx;
-       fr_schedule_receiver_t *sr = arg;
-       fr_schedule_t *sc = sr->sc;
+       fr_schedule_network_t *sn = arg;
+       fr_schedule_t *sc = sn->sc;
        fr_schedule_child_status_t status = FR_CHILD_FAIL;
 
        fr_log(sc->log, L_INFO, "Network starting\n");
 
-       ctx = talloc_init("receiver");
+       ctx = talloc_init("network");
        if (!ctx) {
-               fr_log(sc->log, L_ERR, "Network %d - Failed allocating memory", sr->id);
+               fr_log(sc->log, L_ERR, "Network %d - Failed allocating memory", sn->id);
                goto fail;
        }
 
-       sr->rc = fr_receiver_create(ctx, sc->log, sc->num_transports, sc->transports);
-       if (!sr->rc) {
-               fr_log(sc->log, L_ERR, "Network %d - Failed creating network: %s", sr->id, fr_strerror());
+       sn->rc = fr_network_create(ctx, sc->log, sc->num_transports, sc->transports);
+       if (!sn->rc) {
+               fr_log(sc->log, L_ERR, "Network %d - Failed creating network: %s", sn->id, fr_strerror());
                goto fail;
        }
 
-       sr->status = FR_CHILD_RUNNING;
+       sn->status = FR_CHILD_RUNNING;
 
        /*
         *      Tell the originator that the thread has started.
@@ -333,22 +333,22 @@ static void *fr_schedule_receiver_thread(void *arg)
        /*
         *      Do all of the work.
         */
-       fr_receiver(sr->rc);
+       fr_network(sn->rc);
 
        /*
         *      Talloc ordering issues. We want to be independent of
         *      how talloc walks it's children, and ensure that some
         *      things are freed in a specific order.
         */
-       fr_receiver_destroy(sr->rc);
-       sr->rc = NULL;
+       fr_network_destroy(sn->rc);
+       sn->rc = NULL;
 
        status = FR_CHILD_EXITED;
 
 fail:
        if (ctx) talloc_free(ctx);
 
-       sr->status = status;
+       sn->status = status;
 
        fr_log(sc->log, L_INFO, "Network exiting");
 
@@ -452,20 +452,20 @@ fr_schedule_t *fr_schedule_create(TALLOC_CTX *ctx, fr_log_t *logger, int max_inp
        /*
         *      Create the network thread first.
         */
-       sc->sr = talloc_zero(sc, fr_schedule_receiver_t);
-       sc->sr->sc = sc;
-       sc->sr->id = 0;
+       sc->sn = talloc_zero(sc, fr_schedule_network_t);
+       sc->sn->sc = sc;
+       sc->sn->id = 0;
 
-       rcode = pthread_create(&sc->sr->pthread_id, &attr, fr_schedule_receiver_thread, sc->sr);
+       rcode = pthread_create(&sc->sn->pthread_id, &attr, fr_schedule_network_thread, sc->sn);
        if (rcode != 0) {
                fr_strerror_printf("Failed creating network thread: %s", fr_syserror(errno));
                goto fail;
        }
 
        SEM_WAIT_INTR(&sc->semaphore);
-       if (sc->sr->status != FR_CHILD_RUNNING) {
+       if (sc->sn->status != FR_CHILD_RUNNING) {
        fail:
-               TALLOC_FREE(sc->sr);
+               TALLOC_FREE(sc->sn);
                fr_schedule_destroy(sc);
                return NULL;
        }
@@ -621,8 +621,8 @@ int fr_schedule_destroy(fr_schedule_t *sc)
        /*
         *      If the network thread is running, tell it to exit.
         */
-       if (sc->sr->status == FR_CHILD_RUNNING) {
-               fr_receiver_exit(sc->sr->rc);
+       if (sc->sn->status == FR_CHILD_RUNNING) {
+               fr_network_exit(sc->sn->rc);
                SEM_WAIT_INTR(&sc->semaphore);
        }
 
@@ -649,7 +649,7 @@ int fr_schedule_destroy(fr_schedule_t *sc)
  */
 int fr_schedule_socket_add(fr_schedule_t *sc, int fd, void *ctx, fr_transport_t *transport)
 {
-       return fr_receiver_socket_add(sc->sr->rc, fd, ctx, transport);
+       return fr_network_socket_add(sc->sn->rc, fd, ctx, transport);
 }
 
 
index 2c7fb5bc7e798ccc7cc723477642a909aa3dda6e..7d5b1393f008664db4303df43f6e87692d149980 100644 (file)
@@ -208,7 +208,7 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_
                fr_log(worker->log, L_DBG, "\t%saq noop", worker->name);
                return;
 
-       case FR_CHANNEL_DATA_READY_RECEIVER:
+       case FR_CHANNEL_DATA_READY_NETWORK:
                rad_assert(0 == 1);
                fr_log(worker->log, L_DBG, "\t%saq data ready ? MASTER ?", worker->name);
                break;
index 80fca30e5b4762e35558f67d7611c6ff151ca504..fe8745ea38e32c11a0a0f13233033dcac9d61934 100644 (file)
@@ -229,7 +229,7 @@ check_close:
                        MPRINT1("Master got channel event %d\n", ce);
                        
                        switch (ce) {
-                       case FR_CHANNEL_DATA_READY_RECEIVER:
+                       case FR_CHANNEL_DATA_READY_NETWORK:
                                MPRINT1("Master got data ready signal\n");
                                rad_assert(new_channel == channel);
 
index 5304265d9314df4496d15a49a11fd0de9a6caf12..15d0861a88200799860403c16d92f2d251193f37 100644 (file)
@@ -431,7 +431,7 @@ static void master_process(TALLOC_CTX *ctx)
                        MPRINT1("Master got channel event %d\n", ce);
                        
                        switch (ce) {
-                       case FR_CHANNEL_DATA_READY_RECEIVER:
+                       case FR_CHANNEL_DATA_READY_NETWORK:
                                MPRINT1("Master got data ready signal\n");
 
                                reply = fr_channel_recv_reply(ch);
index 5b764da5db0628620e921b9754a8043871a23e1c..170a1b6bc859b3fe6343d0efa862f94775a25c9c 100644 (file)
@@ -355,7 +355,7 @@ check_close:
                        MPRINT1("Master got channel event %d\n", ce);
                        
                        switch (ce) {
-                       case FR_CHANNEL_DATA_READY_RECEIVER:
+                       case FR_CHANNEL_DATA_READY_NETWORK:
                                MPRINT1("Master got data ready signal\n");
 
                                reply = fr_channel_recv_reply(ch);