Use native threads in host resolver so that it works even if processor has no threads

This commit is contained in:
Tobias Brunner
2012-10-18 12:26:49 +02:00
parent b4f6c39e55
commit d377556863
+49 -17
View File
@@ -20,7 +20,7 @@
#include "host_resolver.h" #include "host_resolver.h"
#include <debug.h> #include <debug.h>
#include <processing/jobs/callback_job.h> #include <library.h>
#include <threading/condvar.h> #include <threading/condvar.h>
#include <threading/mutex.h> #include <threading/mutex.h>
#include <threading/thread.h> #include <threading/thread.h>
@@ -90,6 +90,11 @@ struct private_host_resolver_t {
*/ */
u_int busy_threads; u_int busy_threads;
/**
* Pool of threads, thread_t*
*/
linked_list_t *pool;
/** /**
* TRUE if no new queries are accepted * TRUE if no new queries are accepted
*/ */
@@ -153,31 +158,42 @@ static bool query_equals(query_t *this, query_t *other)
/** /**
* Main function of resolver threads * Main function of resolver threads
*/ */
static job_requeue_t resolve_hosts(private_host_resolver_t *this) static void *resolve_hosts(private_host_resolver_t *this)
{ {
struct addrinfo hints, *result; struct addrinfo hints, *result;
query_t *query; query_t *query;
int error; int error;
bool old, timed_out; bool old, timed_out;
while (TRUE)
{
this->mutex->lock(this->mutex); this->mutex->lock(this->mutex);
while (this->queue->remove_first(this->queue, (void**)&query) != SUCCESS)
{
thread_cleanup_push((thread_cleanup_t)this->mutex->unlock, this->mutex); thread_cleanup_push((thread_cleanup_t)this->mutex->unlock, this->mutex);
while (this->queue->remove_first(this->queue,
(void**)&query) != SUCCESS)
{
old = thread_cancelability(TRUE); old = thread_cancelability(TRUE);
timed_out = this->new_query->timed_wait(this->new_query, this->mutex, timed_out = this->new_query->timed_wait(this->new_query,
NEW_QUERY_WAIT_TIMEOUT * 1000); this->mutex, NEW_QUERY_WAIT_TIMEOUT * 1000);
thread_cancelability(old); thread_cancelability(old);
if (timed_out && (this->threads > this->min_threads)) if (this->disabled)
{ {
this->threads--;
thread_cleanup_pop(TRUE); thread_cleanup_pop(TRUE);
return JOB_REQUEUE_NONE; return NULL;
}
else if (timed_out && (this->threads > this->min_threads))
{ /* terminate this thread by detaching it */
thread_t *thread = thread_current();
this->threads--;
this->pool->remove(this->pool, thread, NULL);
thread_cleanup_pop(TRUE);
thread->detach(thread);
return NULL;
} }
thread_cleanup_pop(FALSE);
} }
this->busy_threads++; this->busy_threads++;
this->mutex->unlock(this->mutex); thread_cleanup_pop(TRUE);
memset(&hints, 0, sizeof(hints)); memset(&hints, 0, sizeof(hints));
hints.ai_family = query->family; hints.ai_family = query->family;
@@ -204,7 +220,8 @@ static job_requeue_t resolve_hosts(private_host_resolver_t *this)
query->done->broadcast(query->done); query->done->broadcast(query->done);
this->mutex->unlock(this->mutex); this->mutex->unlock(this->mutex);
query_destroy(query); query_destroy(query);
return JOB_REQUEUE_DIRECT; }
return NULL;
} }
METHOD(host_resolver_t, resolve, host_t*, METHOD(host_resolver_t, resolve, host_t*,
@@ -250,12 +267,15 @@ METHOD(host_resolver_t, resolve, host_t*,
ref_get(&query->refcount); ref_get(&query->refcount);
if (this->busy_threads == this->threads && if (this->busy_threads == this->threads &&
this->threads < this->max_threads) this->threads < this->max_threads)
{
thread_t *thread;
thread = thread_create((thread_main_t)resolve_hosts, this);
if (thread)
{ {
this->threads++; this->threads++;
lib->processor->queue_job(lib->processor, this->pool->insert_last(this->pool, thread);
(job_t*)callback_job_create_with_prio( }
(callback_job_cb_t)resolve_hosts, this, NULL,
(callback_job_cancel_t)return_false, JOB_PRIO_CRITICAL));
} }
query->done->wait(query->done, this->mutex); query->done->wait(query->done, this->mutex);
this->mutex->unlock(this->mutex); this->mutex->unlock(this->mutex);
@@ -282,13 +302,24 @@ METHOD(host_resolver_t, flush, void,
this->queue->destroy_function(this->queue, (void*)query_destroy); this->queue->destroy_function(this->queue, (void*)query_destroy);
this->queue = linked_list_create(); this->queue = linked_list_create();
this->disabled = TRUE; this->disabled = TRUE;
/* this will already terminate most idle threads */
this->new_query->broadcast(this->new_query);
this->mutex->unlock(this->mutex); this->mutex->unlock(this->mutex);
} }
METHOD(host_resolver_t, destroy, void, METHOD(host_resolver_t, destroy, void,
private_host_resolver_t *this) private_host_resolver_t *this)
{ {
this->queue->destroy_function(this->queue, (void*)query_signal_and_destroy); thread_t *thread;
flush(this);
this->pool->invoke_offset(this->pool, offsetof(thread_t, cancel));
while (this->pool->remove_first(this->pool, (void**)&thread) == SUCCESS)
{
thread->join(thread);
}
this->pool->destroy(this->pool);
this->queue->destroy(this->queue);
this->queries->destroy(this->queries); this->queries->destroy(this->queries);
this->new_query->destroy(this->new_query); this->new_query->destroy(this->new_query);
this->mutex->destroy(this->mutex); this->mutex->destroy(this->mutex);
@@ -311,6 +342,7 @@ host_resolver_t *host_resolver_create()
.queries = hashtable_create((hashtable_hash_t)query_hash, .queries = hashtable_create((hashtable_hash_t)query_hash,
(hashtable_equals_t)query_equals, 8), (hashtable_equals_t)query_equals, 8),
.queue = linked_list_create(), .queue = linked_list_create(),
.pool = linked_list_create(),
.mutex = mutex_create(MUTEX_TYPE_DEFAULT), .mutex = mutex_create(MUTEX_TYPE_DEFAULT),
.new_query = condvar_create(CONDVAR_TYPE_DEFAULT), .new_query = condvar_create(CONDVAR_TYPE_DEFAULT),
); );