* 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
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);
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);
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),
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);
#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;
}
* 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;
}
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;
}
* 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;
}
* 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;
}
*/
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());
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
*/
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
return NULL;
}
+ worker->name = "";
+
worker->channel = talloc_zero_array(worker, fr_channel_t *, max_channels);
if (!worker->channel) {
talloc_free(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);
}
* 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);
}
* 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);
}
}
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);
+}