From: Alan T. DeKok Date: Wed, 5 Apr 2017 21:19:17 +0000 (-0400) Subject: set names for workers X-Git-Url: http://git.ipfire.org/cgi-bin/gitweb.cgi?a=commitdiff_plain;h=8851044b03152145e6cf58c9727d46fb4287ebd7;p=thirdparty%2Ffreeradius-server.git set names for workers --- diff --git a/src/lib/io/schedule.c b/src/lib/io/schedule.c index 44c8f7dee0e..4fb71ceb3e5 100644 --- a/src/lib/io/schedule.c +++ b/src/lib/io/schedule.c @@ -202,6 +202,7 @@ static void *fr_schedule_worker_thread(void *arg) fr_schedule_worker_t *sw = arg; fr_schedule_t *sc = sw->sc; fr_schedule_child_status_t status = FR_CHILD_FAIL; + char buffer[32]; fr_log(sc->log, L_INFO, "Worker %d starting\n", sw->id); @@ -217,6 +218,9 @@ static void *fr_schedule_worker_thread(void *arg) goto fail; } + snprintf(buffer, sizeof(buffer), "thread %d - ", sw->id); + fr_worker_name(sw->worker, buffer); + /* * @todo make this a registry */ diff --git a/src/lib/io/worker.c b/src/lib/io/worker.c index d4a436cab1d..2192f921fc6 100644 --- a/src/lib/io/worker.c +++ b/src/lib/io/worker.c @@ -79,6 +79,8 @@ typedef struct fr_worker_heap_t { * A worker which takes packets from a master, and processes them. */ struct fr_worker_t { + char const *name; //!< name of this worker + int kq; //!< my kq fr_log_t *log; //!< log destination @@ -162,14 +164,14 @@ static void fr_worker_drain_input(fr_worker_t *worker, fr_channel_t *ch, fr_chan if (!cd) { cd = fr_channel_recv_request(ch); if (!cd) { - fr_log(worker->log, L_DBG, "\tno data?"); + fr_log(worker->log, L_DBG, "\t%sno data?", worker->name); return; } } do { worker->num_requests++; - fr_log(worker->log, L_DBG, "\treceived request %d", worker->num_requests); + fr_log(worker->log, L_DBG, "\t%sreceived request %d", worker->name, worker->num_requests); cd->channel.ch = ch; WORKER_HEAP_INSERT(to_decode, cd, request.list); } while ((cd = fr_channel_recv_request(ch)) != NULL); @@ -195,30 +197,30 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_ ce = fr_channel_service_message(now, &ch, data, data_size); switch (ce) { case FR_CHANNEL_ERROR: - fr_log(worker->log, L_DBG, "\taq error"); + fr_log(worker->log, L_DBG, "\t%saq error", worker->name); return; case FR_CHANNEL_EMPTY: - fr_log(worker->log, L_DBG, "\taq empty"); + fr_log(worker->log, L_DBG, "\t%saq empty", worker->name); return; case FR_CHANNEL_NOOP: - fr_log(worker->log, L_DBG, "\taq noop"); + fr_log(worker->log, L_DBG, "\t%saq noop", worker->name); return; case FR_CHANNEL_DATA_READY_RECEIVER: rad_assert(0 == 1); - fr_log(worker->log, L_DBG, "\taq data ready ? MASTER ?"); + fr_log(worker->log, L_DBG, "\t%saq data ready ? MASTER ?", worker->name); break; case FR_CHANNEL_DATA_READY_WORKER: rad_assert(ch != NULL); - fr_log(worker->log, L_DBG, "\taq data ready"); + fr_log(worker->log, L_DBG, "\t%saq data ready", worker->name); fr_worker_drain_input(worker, ch, NULL); break; case FR_CHANNEL_OPEN: - fr_log(worker->log, L_DBG, "\taq channel open"); + fr_log(worker->log, L_DBG, "\t%saq channel open", worker->name); rad_assert(ch != NULL); @@ -229,7 +231,7 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_ if (worker->channel[i] != NULL) continue; worker->channel[i] = ch; - fr_log(worker->log, L_DBG, "\treceived channel %p into array entry %d", ch, i); + fr_log(worker->log, L_DBG, "\t%sreceived channel %p into array entry %d", worker->name, ch, i); ms = fr_message_set_create(worker, worker->message_set_size, sizeof(fr_channel_data_t), @@ -246,7 +248,7 @@ static void fr_worker_channel_callback(void *ctx, void const *data, size_t data_ break; case FR_CHANNEL_CLOSE: - fr_log(worker->log, L_DBG, "\taq channel close"); + fr_log(worker->log, L_DBG, "\t%saq channel close", worker->name); rad_assert(ch != NULL); @@ -301,7 +303,7 @@ static void fr_worker_evfilt_user(UNUSED int kq, struct kevent const *kev, void #endif if (!fr_control_message_service_kevent(worker->control, kev)) { - fr_log(worker->log, L_DBG, "\tkevent not for us!"); + fr_log(worker->log, L_DBG, "\t%skevent not for us!", worker->name); return; } @@ -374,7 +376,7 @@ static void fr_worker_nak(fr_worker_t *worker, fr_channel_data_t *cd, fr_time_t * Send the reply, which also polls the request queue. */ if (fr_channel_send_reply(ch, reply, &cd) < 0) { - fr_log(worker->log, L_DBG, "\tfails sending reply"); + fr_log(worker->log, L_DBG, "\t%sfails sending reply", worker->name); cd = NULL; } @@ -418,7 +420,7 @@ static void fr_worker_send_reply(fr_worker_t *worker, REQUEST *request, size_t s encoded = request->transport->encode(request->packet_ctx, request, reply->m.data, reply->m.rb_size); if (encoded < 0) { - fr_log(worker->log, L_DBG, "\tfails encode"); + fr_log(worker->log, L_DBG, "\t%sfails encode", worker->name); encoded = 0; } @@ -455,7 +457,7 @@ static void fr_worker_send_reply(fr_worker_t *worker, REQUEST *request, size_t s * Send the reply, which also polls the request queue. */ if (fr_channel_send_reply(ch, reply, &cd) < 0) { - fr_log(worker->log, L_DBG, "\tfails sending reply"); + fr_log(worker->log, L_DBG, "\t%sfails sending reply", worker->name); cd = NULL; } @@ -650,7 +652,7 @@ static REQUEST *fr_worker_get_request(fr_worker_t *worker, fr_time_t now) * in the queue. Delete it, and go get another one. */ if (cd->request.start_time && (cd->m.when != *cd->request.start_time)) { - fr_log(worker->log, L_DBG, "\tIGNORING old message"); + fr_log(worker->log, L_DBG, "\t%sIGNORING old message", worker->name); fr_worker_nak(worker, cd, fr_time()); cd = NULL; } @@ -708,7 +710,7 @@ static REQUEST *fr_worker_get_request(fr_worker_t *worker, fr_time_t now) */ rcode = worker->transports[cd->transport]->decode(cd->packet_ctx, cd->m.data, cd->m.data_size, request); if (rcode < 0) { - fr_log(worker->log, L_DBG, "\tFAILED decode of request %zd", request->number); + fr_log(worker->log, L_DBG, "\t%sFAILED decode of request %zd", worker->name, request->number); talloc_free(ctx); nak: fr_worker_nak(worker, cd, fr_time()); @@ -760,7 +762,7 @@ static void fr_worker_run_request(fr_worker_t *worker, REQUEST *request) ssize_t size = 0; fr_transport_final_t final; - fr_log(worker->log, L_DBG, "(%zd) running", request->number); + fr_log(worker->log, L_DBG, "\t%s running request (%zd)", worker->name, request->number); /* * If we still have the same packet, and the channel is @@ -852,12 +854,13 @@ static int fr_worker_idle(void *ctx, struct timeval *wake) */ if (!sleeping) return 1; - fr_log(worker->log, L_DBG, "\tsleeping running %zd, localized %zd, to_decode %zd", + fr_log(worker->log, L_DBG, "\t%ssleeping running %zd, localized %zd, to_decode %zd", + worker->name, fr_heap_num_elements(worker->runnable), fr_heap_num_elements(worker->localized.heap), fr_heap_num_elements(worker->to_decode.heap)); - fr_log(worker->log, L_DBG, "\trequests %d, decoded %d, replied %d", - worker->num_requests, worker->num_decoded, worker->num_replies); + fr_log(worker->log, L_DBG, "\t%srequests %d, decoded %d, replied %d", + worker->name, worker->num_requests, worker->num_decoded, worker->num_replies); /* * Nothing more to do, and the event loop has us sleeping @@ -976,6 +979,8 @@ nomem: return NULL; } + worker->name = ""; + worker->channel = talloc_zero_array(worker, fr_channel_t *, max_channels); if (!worker->channel) { talloc_free(worker); @@ -1097,21 +1102,24 @@ void fr_worker(fr_worker_t *worker) * the event loop, but we don't wait for events. */ wait_for_event = (fr_heap_num_elements(worker->runnable) == 0); - fr_log(worker->log, L_DBG, "\tWaiting for events %d", wait_for_event); + fr_log(worker->log, L_DBG, "\t%sWaiting for events %d", worker->name, 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(worker->el, wait_for_event); - fr_log(worker->log, L_DBG, "\tGot num_events %d", num_events); - if (num_events < 0) break; + fr_log(worker->log, L_DBG, "\t%sGot num_events %d", worker->name, num_events); + if (num_events < 0) { + fr_log(worker->log, L_ERR, "Failed corraling events: %s", fr_strerror()); + break; + } /* * Service outstanding events. */ if (num_events > 0) { - fr_log(worker->log, L_DBG, "\tservicing events"); + fr_log(worker->log, L_DBG, "\t%sservicing events", worker->name); fr_event_service(worker->el); } @@ -1121,7 +1129,7 @@ void fr_worker(fr_worker_t *worker) * Ten times a second, check for timeouts on incoming packets. */ if ((now - worker->checked_timeout) > (NANOSEC / 10)) { - fr_log(worker->log, L_DBG, "\tchecking timeouts"); + fr_log(worker->log, L_DBG, "\t%schecking timeouts", worker->name); fr_worker_check_timeouts(worker, now); } @@ -1135,7 +1143,7 @@ void fr_worker(fr_worker_t *worker) * Run the request, and either track it as * yielded, or send a reply. */ - fr_log(worker->log, L_DBG, "\trunning request %p", request); + fr_log(worker->log, L_DBG, "\t%srunning request (%zd)", worker->name, request->number); fr_worker_run_request(worker, request); } } @@ -1211,3 +1219,21 @@ fr_channel_t *fr_worker_channel_create(fr_worker_t const *worker, TALLOC_CTX *ct return ch; } + + +/** Set the name of a worker. + * + * Called by the master (i.e. network) thread when it needs to create + * a new channel to a particuler worker. + * + * @param[in] worker the worker + * @param[in] name the name to set for the worker. (strdup'd by the worker) + */ +void fr_worker_name(fr_worker_t *worker, char const *name) +{ +#ifndef NDEBUG + talloc_get_type_abort(worker, fr_worker_t); +#endif + + worker->name = talloc_strdup(worker, name); +} diff --git a/src/lib/io/worker.h b/src/lib/io/worker.h index 706817c5aa5..0476e44ae9a 100644 --- a/src/lib/io/worker.h +++ b/src/lib/io/worker.h @@ -50,6 +50,7 @@ int fr_worker_kq(fr_worker_t *worker) CC_HINT(nonnull); void fr_worker(fr_worker_t *worker) CC_HINT(nonnull); void fr_worker_exit(fr_worker_t *worker) CC_HINT(nonnull); void fr_worker_debug(fr_worker_t *worker, FILE *fp) CC_HINT(nonnull); +void fr_worker_name(fr_worker_t *worker, char const *name) CC_HINT(nonnull); fr_channel_t *fr_worker_channel_create(fr_worker_t const *worker, TALLOC_CTX *ctx, fr_control_t *master) CC_HINT(nonnull); #ifdef __cplusplus