From: Alan T. DeKok Date: Mon, 24 Apr 2017 15:02:09 +0000 (-0400) Subject: renamed receiver to network. X-Git-Url: http://git.ipfire.org/cgi-bin/gitweb.cgi?a=commitdiff_plain;h=e34d37de7f7ee320ef3fff97d6a01c75b8e1ff32;p=thirdparty%2Ffreeradius-server.git renamed receiver to network. It's not perfect, but it's clearer and makes more sense --- diff --git a/src/lib/io/all.mk b/src/lib/io/all.mk index cfb22b5786e..4d986c6b5e4 100644 --- a/src/lib/io/all.mk +++ b/src/lib/io/all.mk @@ -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) diff --git a/src/lib/io/channel.c b/src/lib/io/channel.c index c25ff3545c8..0c57c1e1d53 100644 --- a/src/lib/io/channel.c +++ b/src/lib/io/channel.c @@ -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: diff --git a/src/lib/io/channel.h b/src/lib/io/channel.h index 880ed238f8c..7585140e0d9 100644 --- a/src/lib/io/channel.h +++ b/src/lib/io/channel.h @@ -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, diff --git a/src/lib/io/receiver.c b/src/lib/io/network.c similarity index 63% rename from src/lib/io/receiver.c rename to src/lib/io/network.c index 84276be3730..7dd761dfd87 100644 --- a/src/lib/io/receiver.c +++ b/src/lib/io/network.c @@ -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 */ @@ -31,20 +31,20 @@ RCSID("$Id$") #include #include #include -#include +#include #include -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)); } diff --git a/src/lib/io/receiver.h b/src/lib/io/network.h similarity index 62% rename from src/lib/io/receiver.h rename to src/lib/io/network.h index dcc655646f2..002d88ed48b 100644 --- a/src/lib/io/receiver.h +++ b/src/lib/io/network.h @@ -18,12 +18,12 @@ /** * $Id$ * - * @file io/receiver.h + * @file io/network.h * @brief Receive packets * * @copyright 2016 Alan DeKok */ -RCSIDH(receiver_h, "$Id$") +RCSIDH(network_h, "$Id$") #include @@ -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 */ diff --git a/src/lib/io/schedule.c b/src/lib/io/schedule.c index 4fb71ceb3e5..fcba83344a9 100644 --- a/src/lib/io/schedule.c +++ b/src/lib/io/schedule.c @@ -30,7 +30,7 @@ RCSID("$Id$") #include #include -#include +#include #ifdef HAVE_PTHREAD_H #include @@ -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); } diff --git a/src/lib/io/worker.c b/src/lib/io/worker.c index 2c7fb5bc7e7..7d5b1393f00 100644 --- a/src/lib/io/worker.c +++ b/src/lib/io/worker.c @@ -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; diff --git a/src/tests/util/channel_test.c b/src/tests/util/channel_test.c index 80fca30e5b4..fe8745ea38e 100644 --- a/src/tests/util/channel_test.c +++ b/src/tests/util/channel_test.c @@ -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); diff --git a/src/tests/util/radius1_test.c b/src/tests/util/radius1_test.c index 5304265d931..15d0861a882 100644 --- a/src/tests/util/radius1_test.c +++ b/src/tests/util/radius1_test.c @@ -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); diff --git a/src/tests/util/worker_test.c b/src/tests/util/worker_test.c index 5b764da5db0..170a1b6bc85 100644 --- a/src/tests/util/worker_test.c +++ b/src/tests/util/worker_test.c @@ -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);