Files
unbound/services/authload.c

922 lines
24 KiB
C

/*
* services/authload.c - authoritative zone load thread
*
* Copyright (c) 2026, NLnet Labs. All rights reserved.
*
* This software is open source.
*
* Redistribution and use in source and binary forms, with or without
* modification, are permitted provided that the following conditions
* are met:
*
* Redistributions of source code must retain the above copyright notice,
* this list of conditions and the following disclaimer.
*
* Redistributions in binary form must reproduce the above copyright notice,
* this list of conditions and the following disclaimer in the documentation
* and/or other materials provided with the distribution.
*
* Neither the name of the NLNET LABS nor the names of its contributors may
* be used to endorse or promote products derived from this software without
* specific prior written permission.
*
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
* A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT
* HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,
* SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED
* TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
* PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
* LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
* NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
* SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
*/
/**
* \file
*
* This file contains the auth load thread. This loads authority zone
* and RPZ zone information in a thread, in a separate memory structure.
* When it is done, the information is swapped over to the running server.
*/
#include "config.h"
#include "services/authload.h"
#include "daemon/worker.h"
#include "daemon/daemon.h"
#include "services/authzone.h"
#include "libunbound/authload.h"
#include "util/net_help.h"
#include "util/log.h"
#include "util/ub_event.h"
#include "util/timeval_func.h"
#include "util/data/dname.h"
/** Get memory use of buffer. */
static size_t
buffer_get_mem(struct sldns_buffer* buf)
{
if(!buf) return 0;
return sizeof(*buf) + (buf->_data?buf->_capacity:0);
}
/** Auth load notification to string, for descriptive purposes. */
static const char*
auth_load_notification_to_string(enum auth_load_notification_type status)
{
switch(status) {
case auth_load_notification_exit:
return "auth_load_notification_exit";
default:
break;
}
return "unknown_auth_load_notification_value";
}
/** delete chunks */
static void
auth_chunk_list_delete(struct auth_chunk* first)
{
struct auth_chunk* c = first, *cn;
while(c) {
cn = c->next;
free(c->data);
free(c);
c = cn;
}
}
/** Delete auth load task item */
static void
auth_load_task_delete(struct auth_load_task* task)
{
if(!task)
return;
free(task->name);
free(task->host);
free(task->file);
auth_chunk_list_delete(task->chunks_first);
free(task);
}
/** Create new auth load task item */
static struct auth_load_task*
auth_load_task_create(void)
{
struct auth_load_task* task = (struct auth_load_task*)calloc(1,
sizeof(*task));
return task;
}
/** Pick up the work content of task transfer of auth xfr */
static int
auth_load_task_pickup_xfr(struct auth_load_task* task, struct auth_xfer* xfr)
{
task->name = memdup(xfr->name, xfr->namelen);
if(!task->name)
return 0;
task->namelen = xfr->namelen;
task->dclass = xfr->dclass;
if(xfr->task_transfer->master && xfr->task_transfer->master->host) {
task->host = strdup(xfr->task_transfer->master->host);
if(!task->host)
return 0;
}
if(xfr->task_transfer->master && xfr->task_transfer->master->file) {
task->file = strdup(xfr->task_transfer->master->file);
if(!task->file)
return 0;
}
if(xfr->task_transfer->master)
task->on_http = xfr->task_transfer->master->http;
task->on_ixfr = xfr->task_transfer->on_ixfr;
task->on_ixfr_is_axfr = xfr->task_transfer->on_ixfr_is_axfr;
task->serial = xfr->serial;
if(xfr->task_transfer->chunks_first) {
task->chunks_first = xfr->task_transfer->chunks_first;
task->chunks_last = xfr->task_transfer->chunks_last;
task->chunks_total = xfr->task_transfer->chunks_total;
/* The task now has the chunks. Remove them from the
* xfr structure. */
xfr->task_transfer->chunks_first = 0;
xfr->task_transfer->chunks_last = 0;
xfr->task_transfer->chunks_total = 0;
}
if(task->on_http)
task->task_type = AUTH_LOAD_TASK_HTTPCHUNKS;
else task->task_type = AUTH_LOAD_TASK_TRANSFER;
return 1;
}
/** Create xfr task */
static struct auth_load_task*
auth_load_task_create_xfr(struct auth_xfer* xfr, struct worker* worker)
{
struct auth_load_task* task = auth_load_task_create();
if(!task) {
log_err("out of memory");
return 0;
}
task->worker = worker;
if(!auth_load_task_pickup_xfr(task, xfr)) {
log_err("out of memory");
auth_load_task_delete(task);
return 0;
}
return task;
}
int
auth_load_thread_poll_for_quit(struct auth_load_thread* thr)
{
int inevent, loopexit = 0;
uint8_t cmd;
ssize_t ret;
if(!thr)
return 0;
if(thr->need_to_quit)
return 1;
/* Is there data? */
if(!sock_poll_timeout(thr->commpair[1], 0, 1, 0, &inevent)) {
log_err("auth_load_thread_poll_for_quit: poll failed");
return 0;
}
if(!inevent)
return 0;
/* Read the data */
while(1) {
if(++loopexit > 200) {
log_err("auth_load_thread_poll_for_quit: recv loops %s",
sock_strerror(errno));
return 0;
}
ret = recv(thr->commpair[1], ((char*)&cmd), sizeof(cmd), 0);
if(ret == -1) {
if(
#ifndef USE_WINSOCK
errno == EINTR || errno == EAGAIN
# ifdef EWOULDBLOCK
|| errno == EWOULDBLOCK
# endif
#else
WSAGetLastError() == WSAEINTR ||
WSAGetLastError() == WSAEINPROGRESS ||
WSAGetLastError() == WSAEWOULDBLOCK
#endif
)
continue; /* Try again. */
log_err("auth_load_thread_poll_for_quit: recv: %s",
sock_strerror(errno));
return 0;
} else if(ret == 0) {
log_err("auth_load_thread_poll_for_quit: recv: EOF");
return 0;
}
break;
}
if(cmd == auth_load_notification_exit) {
thr->need_to_quit = 1;
verbose(VERB_ALGO, "auth load: exit notification received");
return 1;
}
log_err("auth_load_thread_poll_for_quit: unknown notification status "
"received: %d %s", cmd, auth_load_notification_to_string(cmd));
return 0;
}
/** Signal the worker connected to an auth load thread the status */
static void
auth_load_thread_signal_worker(struct auth_load_thread* thr, int status)
{
int outevent, loopexit = 0;
ssize_t ret;
uint8_t to_send;
verbose(VERB_ALGO, "auth load thread: send status %d", status);
/* Make a blocking attempt to send. But meanwhile stay responsive,
* once in a while for quit commands. In case the server has to quit. */
/* see if there is incoming quit signals */
if(auth_load_thread_poll_for_quit(thr))
return;
to_send = (uint8_t)status;
while(1) {
if(++loopexit > 200) {
log_err("auth load thread: could not send status");
return;
}
/* wait for socket to become writable */
if(!sock_poll_timeout(thr->commpair[1],
200, /* msec wait before check for quit, and loop to
wait again. */
0, 1, &outevent)) {
log_err("auth load thread: poll failed");
return;
}
if(auth_load_thread_poll_for_quit(thr))
return;
if(!outevent)
continue;
ret = send(thr->commpair[1], &to_send, 1, 0);
if(ret == -1) {
if(
#ifndef USE_WINSOCK
errno == EINTR || errno == EAGAIN
# ifdef EWOULDBLOCK
|| errno == EWOULDBLOCK
# endif
#else
WSAGetLastError() == WSAEINTR ||
WSAGetLastError() == WSAEINPROGRESS ||
WSAGetLastError() == WSAEWOULDBLOCK
#endif
)
continue; /* Try again. */
log_err("auth load thread signal worker: send: %s",
sock_strerror(errno));
return;
} else if(ret < 1) {
continue;
}
break;
}
}
/** Create proxy auth zone structure, that is used to hold the data
* that is processed. */
static struct auth_zone*
auth_zone_create_proxy(uint8_t* nm, size_t nmlen, uint16_t dclass)
{
struct auth_zone* z = (struct auth_zone*)calloc(1, sizeof(*z));
if(!z) {
return NULL;
}
z->node.key = z;
z->dclass = dclass;
z->namelen = nmlen;
z->namelabs = dname_count_labels(nm);
z->name = memdup(nm, nmlen);
if(!z->name) {
free(z);
return NULL;
}
rbtree_init(&z->data, &auth_data_cmp);
return z;
}
/** Delete proxy auth zone structure */
static void
auth_zone_delete_proxy(struct auth_zone* z)
{
if(!z)
return;
traverse_postorder(&z->data, auth_data_del, NULL);
if(z->rpz)
rpz_delete(z->rpz);
free(z->name);
free(z);
}
/** Calculate memory use of the authload thread for this task.
* The size of the task struct, with the data chunks, and the proxy auth zone
* structure that is created while the other auth zone is used for queries,
* and other added memory.
*/
static void
auth_load_calc_mem(struct auth_load_task* task, struct auth_zone* z,
size_t other)
{
size_t m = 0;
if(verbosity < 8) {
task->mem_used = 0;
return;
}
m += other;
m += sizeof(*task);
m += task->namelen;
m += getmem_str(task->host);
m += getmem_str(task->file);
m += task->chunks_total;
m += auth_zone_get_mem(z);
task->mem_used = m;
}
/** Swap the final zone contents with the live zone */
static void
auth_load_swap_zone(struct auth_load_thread* thr, struct auth_zone* proxyz)
{
rbtree_type data;
struct rpz* rpz;
struct auth_zone* z;
lock_rw_rdlock(&thr->task->worker->env.auth_zones->lock);
z = auth_zone_find(thr->task->worker->env.auth_zones,
thr->task->name, thr->task->namelen, thr->task->dclass);
if(!z) {
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
verbose(VERB_ALGO, "auth zone missing after auth load.");
return;
}
lock_rw_wrlock(&z->lock);
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
data = proxyz->data;
proxyz->data = z->data;
z->data = data;
rpz = proxyz->rpz;
proxyz->rpz = z->rpz;
z->rpz = rpz;
lock_rw_unlock(&z->lock);
}
/** Process http transfer */
static int
auth_load_process_http(struct auth_load_thread* thr)
{
struct auth_load_task* task = thr->task;
struct sldns_buffer* scratch_buffer;
struct auth_zone* z;
size_t scratch_mem;
scratch_buffer = sldns_buffer_new(sldns_buffer_capacity(
thr->task->worker->env.scratch_buffer));
if(!scratch_buffer) {
log_err("out of memory");
return 0;
}
scratch_mem = buffer_get_mem(scratch_buffer);
z = auth_zone_create_proxy(task->name, task->namelen, task->dclass);
if(!z) {
log_err("out of memory");
sldns_buffer_free(scratch_buffer);
return 0;
}
if(auth_load_thread_poll_for_quit(thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
xfr_http_preview(task->file, task->chunks_first);
if(!xfr_http_syntax_check(task->name, task->namelen, task->dclass,
task->host, task->file, task->chunks_first, scratch_buffer)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
if(auth_load_thread_poll_for_quit(thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
if(!xfr_apply_http(task->name, task->namelen, task->host, task->file,
task->chunks_first, z, scratch_buffer, thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
sldns_buffer_free(scratch_buffer);
if(z->rpz)
rpz_finish_config(z->rpz);
if(auth_load_thread_poll_for_quit(thr)) {
auth_zone_delete_proxy(z);
return 0;
}
auth_load_calc_mem(task, z, scratch_mem);
auth_load_swap_zone(thr, z);
auth_zone_delete_proxy(z);
return 1;
}
/** Copy RRset and append it to the domain, update last pointer. */
static int
rrset_append_copy(struct auth_data* domain, struct auth_rrset* rrset,
struct auth_rrset** last)
{
struct auth_rrset* s = calloc(1, sizeof(*s));
if(!s)
return 0;
s->type = rrset->type;
s->data = (struct packed_rrset_data*)memdup(rrset->data,
packed_rrset_sizeof(rrset->data));
if(!s->data) {
free(s);
return 0;
}
packed_rrset_ptr_fixup(s->data);
if(!*last)
domain->rrsets = s;
else (*last)->next = s;
*last = s;
return 1;
}
/** Copy the existing zone for modification */
static int
auth_load_copy_into_zone(struct auth_load_thread* thr, struct auth_zone* proxyz)
{
int count = 0;
struct auth_zone* z;
struct auth_data* d;
lock_rw_rdlock(&thr->task->worker->env.auth_zones->lock);
z = auth_zone_find(thr->task->worker->env.auth_zones,
thr->task->name, thr->task->namelen, thr->task->dclass);
if(!z) {
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
verbose(VERB_ALGO, "auth zone missing for copy for IXFR.");
return 0;
}
lock_rw_rdlock(&z->lock);
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
/* Copy from z into proxyz. */
RBTREE_FOR(d, struct auth_data*, &z->data) {
struct auth_rrset* rrset, *last = NULL;
struct auth_data* proxy_d = az_domain_create(proxyz,
d->name, d->namelen);
if(!proxy_d) {
log_err("out of memory");
lock_rw_unlock(&z->lock);
return 0;
}
for(rrset = d->rrsets; rrset; rrset=rrset->next) {
if(!rrset_append_copy(proxy_d, rrset, &last)) {
log_err("out of memory");
lock_rw_unlock(&z->lock);
return 0;
}
if((count++)%10000 == 0) {
if(auth_load_thread_poll_for_quit(thr)) {
lock_rw_unlock(&z->lock);
return 0;
}
}
}
if((count++)%10000 == 0) {
if(auth_load_thread_poll_for_quit(thr)) {
lock_rw_unlock(&z->lock);
return 0;
}
}
}
lock_rw_unlock(&z->lock);
return 1;
}
/** Process ixfr transfer */
static int
auth_load_process_ixfr(struct auth_load_thread* thr)
{
struct auth_load_task* task = thr->task;
struct sldns_buffer* scratch_buffer;
struct auth_zone* z;
size_t scratch_mem;
scratch_buffer = sldns_buffer_new(sldns_buffer_capacity(
thr->task->worker->env.scratch_buffer));
if(!scratch_buffer) {
log_err("out of memory");
return 0;
}
scratch_mem = buffer_get_mem(scratch_buffer);
z = auth_zone_create_proxy(task->name, task->namelen, task->dclass);
if(!z) {
log_err("out of memory");
sldns_buffer_free(scratch_buffer);
return 0;
}
if(auth_load_thread_poll_for_quit(thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
/* Copy the existing zone for modification, that uses a read lock.
* That then does not interrupt the service of threads. */
if(!auth_load_copy_into_zone(thr, z)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
if(!xfr_apply_ixfr(task->chunks_first, task->serial, z,
scratch_buffer, thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
sldns_buffer_free(scratch_buffer);
if(auth_load_thread_poll_for_quit(thr)) {
auth_zone_delete_proxy(z);
return 0;
}
auth_load_calc_mem(task, z, scratch_mem);
auth_load_swap_zone(thr, z);
auth_zone_delete_proxy(z);
return 1;
}
/** Process axfr transfer */
static int
auth_load_process_axfr(struct auth_load_thread* thr)
{
struct auth_load_task* task = thr->task;
struct sldns_buffer* scratch_buffer;
struct auth_zone* z;
size_t scratch_mem;
scratch_buffer = sldns_buffer_new(sldns_buffer_capacity(
thr->task->worker->env.scratch_buffer));
if(!scratch_buffer) {
log_err("out of memory");
return 0;
}
scratch_mem = buffer_get_mem(scratch_buffer);
z = auth_zone_create_proxy(task->name, task->namelen, task->dclass);
if(!z) {
log_err("out of memory");
sldns_buffer_free(scratch_buffer);
return 0;
}
if(auth_load_thread_poll_for_quit(thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
if(!xfr_apply_axfr(task->chunks_first, z, scratch_buffer, thr)) {
sldns_buffer_free(scratch_buffer);
auth_zone_delete_proxy(z);
return 0;
}
sldns_buffer_free(scratch_buffer);
if(auth_load_thread_poll_for_quit(thr)) {
auth_zone_delete_proxy(z);
return 0;
}
auth_load_calc_mem(task, z, scratch_mem);
auth_load_swap_zone(thr, z);
auth_zone_delete_proxy(z);
return 1;
}
/** In the auth load thread, process the task */
static int
auth_load_thread_process(struct auth_load_thread* thr)
{
struct auth_load_task* task = thr->task;
struct timeval start, end;
if(gettimeofday(&start, NULL) < 0)
log_err("gettimeofday: %s", strerror(errno));
/* apply data */
if(task->on_http) {
if(!auth_load_process_http(thr))
return 0;
} else if(task->on_ixfr && !task->on_ixfr_is_axfr) {
if(!auth_load_process_ixfr(thr))
return 0;
} else {
if(!auth_load_process_axfr(thr))
return 0;
}
if(gettimeofday(&end, NULL) < 0)
log_err("gettimeofday: %s", strerror(errno));
timeval_subtract(&thr->task->time_taken, &end, &start);
return 1;
}
/** The auth load thread. The thread main function. */
static void*
auth_load_thread_main(void* arg)
{
struct auth_load_thread* thr = (struct auth_load_thread*)arg;
int s;
const char name[16] = "unbound/authld"; /* seems to be the safest size
between different OSes */
#if defined(HAVE_GETTID) && !defined(THREADS_DISABLED)
thr->thread_tid = gettid();
if(thr->thread_tid_log)
log_thread_set(&thr->thread_tid);
else
#endif
log_thread_set(&thr->threadnum);
ub_thread_setname(ub_thread_self(), name);
(void)name; /* When setname is not defined, ignore the name variable. */
verbose(VERB_ALGO, "start auth load thread");
s = auth_load_thread_process(thr);
/* The result is sent to the worker, that reaps the thread. */
auth_load_thread_signal_worker(thr, s);
verbose(VERB_ALGO, "stop auth load thread");
return NULL;
}
/** Delete auth load thread structure */
static void
auth_load_thread_delete(struct auth_load_thread* thr)
{
if(!thr)
return;
if(thr->service_event && thr->service_event_is_added) {
ub_event_del(thr->service_event);
thr->service_event_is_added = 0;
}
if(thr->service_event)
ub_event_free(thr->service_event);
if(thr->commpair[0] != -1)
sock_close(thr->commpair[0]);
if(thr->commpair[1] != -1)
sock_close(thr->commpair[1]);
auth_load_task_delete(thr->task);
free(thr);
}
/** Create auth load thread structure */
static struct auth_load_thread*
auth_load_thread_create(struct auth_load_task* task)
{
int numworkers;
struct auth_load_thread* thr = (struct auth_load_thread*)calloc(1,
sizeof(*thr));
if(!thr)
return NULL;
numworkers = task->worker->daemon->num;
/* This number is printed into the logs */
thr->threadnum = numworkers+3;
thr->task = task;
thr->commpair[0] = -1;
thr->commpair[1] = -1;
if(!create_socketpair(thr->commpair, task->worker->daemon->rand)) {
auth_load_thread_delete(thr);
return NULL;
}
#ifdef HAVE_GETTID
thr->thread_tid_log = task->worker->env.cfg->log_thread_id;
#endif
return thr;
}
/** The worker routine that services the auth load connection. */
void
worker_auth_load_service_cb(int ATTR_UNUSED(fd), short ATTR_UNUSED(bits),
void* arg)
{
struct auth_load_thread* thr = (struct auth_load_thread*)arg;
uint8_t recv_item;
ssize_t ret;
struct auth_xfer* xfr;
struct auth_chunk* chunk_list;
struct module_env* env = &thr->task->worker->env;
int ixfr_fail;
struct timeval time_taken;
size_t mem_used, chunks_total;
log_assert(thr->commpair[0] >= 0);
ret = recv(thr->commpair[0], &recv_item, 1, 0);
if(ret == -1) {
if(
#ifndef USE_WINSOCK
errno == EINTR || errno == EAGAIN
# ifdef EWOULDBLOCK
|| errno == EWOULDBLOCK
# endif
#else
WSAGetLastError() == WSAEINTR ||
WSAGetLastError() == WSAEINPROGRESS
#endif
)
return; /* Continue later. */
#ifdef USE_WINSOCK
if(WSAGetLastError() == WSAEWOULDBLOCK) {
ub_winsock_tcp_wouldblock(thr->service_event,
UB_EV_READ);
return; /* Continue later. */
}
#endif
log_err("read status from auth load thread, recv: %s",
sock_strerror(errno));
return;
} else if(ret == 0) {
verbose(VERB_ALGO, "closed connection from auth load thread");
/* handle this like an error */
recv_item = 0;
/* ret<1: No short read on 1 byte, to continue later on */
}
/* Deal with the result of auth load thread */
verbose(VERB_ALGO, "auth load status is %d", (int)recv_item);
verbose(VERB_ALGO, "join with auth load thread");
ub_thread_join(thr->tid);
verbose(VERB_ALGO, "joined with auth load thread");
lock_rw_rdlock(&thr->task->worker->env.auth_zones->lock);
xfr = auth_xfer_find(thr->task->worker->env.auth_zones,
thr->task->name, thr->task->namelen, thr->task->dclass);
if(!xfr) {
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
verbose(VERB_ALGO, "auth load: xfr is gone");
auth_load_thread_delete(thr);
auth_load_info_release_thread(env);
return;
}
lock_basic_lock(&xfr->lock);
lock_rw_unlock(&thr->task->worker->env.auth_zones->lock);
ixfr_fail = thr->task->ixfr_fail;
time_taken = thr->task->time_taken;
mem_used = thr->task->mem_used;
chunks_total = thr->task->chunks_total;
if(thr->task->on_http) {
chunk_list = thr->task->chunks_first;
thr->task->chunks_first = NULL;
thr->task->chunks_last = NULL;
thr->task->chunks_total = 0;
} else {
chunk_list = NULL;
}
auth_load_thread_delete(thr);
auth_load_info_release_thread(env);
xfr_process_load_end_transfer(xfr, env, recv_item, ixfr_fail,
&time_taken, mem_used, chunks_total, chunk_list);
}
/** Attach worker to the auth load thread. */
static int
auth_load_thread_attach(struct auth_load_thread* thr, struct worker* worker)
{
/* Setup listener in worker, that connects via a pipe to the
* auth load thread.
* The listener has to be nonblocking, so the the remote servicing
* thread can continue to service DNS queries.
* The commpair[1] element can stay blocking, it is used by the
* auth load thread. The thread needs to wait at these times, when
* it has to check briefly it can use poll. */
verbose(VERB_ALGO, "auth_load_thread_attach");
fd_set_nonblock(thr->commpair[0]);
if(!comm_base_internal(worker->base)) {
verbose(VERB_ALGO, "auth load thread: no event base");
return 0;
}
thr->service_event = ub_event_new(comm_base_internal(worker->base),
thr->commpair[0], UB_EV_READ | UB_EV_PERSIST,
worker_auth_load_service_cb, thr);
if(!thr->service_event) {
log_err("out of memory");
return 0;
}
if(ub_event_add(thr->service_event, NULL) != 0) {
log_err("out of memory");
return 0;
}
thr->service_event_is_added = 1;
return 1;
}
/** Create and start the auth load thread, with the task */
static int
auth_load_start_thread(struct auth_load_task* task)
{
struct auth_load_thread* thr = auth_load_thread_create(task);
if(!thr) {
log_err("out of memory");
auth_load_task_delete(task);
return 0;
}
if(!auth_load_thread_attach(thr, task->worker)) {
log_err("out of memory");
auth_load_thread_delete(thr);
return 0;
}
/* Start auth load thread */
ub_thread_create(&thr->tid, auth_load_thread_main, thr);
return 1;
}
int auth_load_add_task_xfr(struct auth_xfer* xfr, struct worker* worker)
{
struct auth_load_task* task;
int can_run = 0;
verbose(VERB_ALGO, "auth load add task");
/* Check auth load count */
can_run = 1;
/* Create new thread */
task = auth_load_task_create_xfr(xfr, worker);
if(!task)
return 0;
if(can_run) {
verbose(VERB_ALGO, "auth load start thread");
if(!auth_load_start_thread(task))
return 0;
verbose(VERB_ALGO, "auth load thread started");
return 1;
}
/* Make wait item */
return 0;
}
struct auth_load_general_info* auth_load_info_create(void)
{
struct auth_load_general_info* auth_load_info =
(struct auth_load_general_info*)calloc(1,
sizeof(*auth_load_info));
if(!auth_load_info) {
log_err("malloc failure");
return NULL;
}
lock_basic_init(&auth_load_info->lock);
lock_protect(&auth_load_info->lock,
&auth_load_info->num_auth_load_threads,
sizeof(auth_load_info->num_auth_load_threads));
return auth_load_info;
}
void auth_load_info_delete(struct auth_load_general_info* auth_load_info)
{
if(!auth_load_info)
return;
lock_basic_destroy(&auth_load_info->lock);
free(auth_load_info);
}
int auth_load_info_grab_thread(struct module_env* env)
{
struct auth_load_general_info* auth_load_info =
env->worker->daemon->auth_load_info;
struct config_file* cfg = env->cfg;
int ret = 0;
lock_basic_lock(&auth_load_info->lock);
if(auth_load_info->num_auth_load_threads < cfg->auth_task_threads) {
ret = 1;
auth_load_info->num_auth_load_threads++;
}
lock_basic_unlock(&auth_load_info->lock);
return ret;
}
void auth_load_info_release_thread(struct module_env* env)
{
struct auth_load_general_info* auth_load_info =
env->worker->daemon->auth_load_info;
lock_basic_lock(&auth_load_info->lock);
if(auth_load_info->num_auth_load_threads == 0) {
verbose(VERB_ALGO, "release of auth load thread, but "
"num_auth_load_threads not > 0.");
} else {
auth_load_info->num_auth_load_threads--;
}
lock_basic_unlock(&auth_load_info->lock);
}