asyn-thrdd: more simplifications

- use wakeup sockets non-locked.
- send wakeup notify only in normal control flow (not cancel). close
  wakeup sockets in unlink only.
- remove 5ms thread lifetime wait crutch before pthread_cancel().

Closes #18380
This commit is contained in:
Stefan Eissing 2025-08-23 14:15:13 +02:00 committed by Daniel Stenberg
parent 41923af5fc
commit d57cfc1a8d
No known key found for this signature in database
GPG key ID: 5CC908FDB71E12C2
2 changed files with 62 additions and 68 deletions

View file

@ -65,7 +65,6 @@
#include "curl_threads.h"
#include "select.h"
#include "strdup.h"
#include "curlx/wait.h"
#ifdef USE_ARES
#include <ares.h>
@ -134,41 +133,7 @@ static void addr_ctx_unlink(struct async_thrdd_addr_ctx **paddr_ctx,
DEBUGASSERT(addr_ctx->ref_count);
--addr_ctx->ref_count;
destroy = (!addr_ctx->ref_count && addr_ctx->thrd_done);
#ifndef CURL_DISABLE_SOCKETPAIR
if(!destroy) {
if(!data) { /* Called from thread, transfer still waiting on results. */
if(addr_ctx->sock_pair[1] != CURL_SOCKET_BAD) {
#ifdef USE_EVENTFD
const uint64_t buf[1] = { 1 };
#else
const char buf[1] = { 1 };
#endif
/* Thread is done, notify transfer */
if(wakeup_write(addr_ctx->sock_pair[1], buf, sizeof(buf)) < 0) {
/* update sock_error to errno */
addr_ctx->sock_error = SOCKERRNO;
}
}
}
else { /* transfer going away, thread still running */
#ifndef USE_EVENTFD
if(addr_ctx->sock_pair[1] != CURL_SOCKET_BAD) {
wakeup_close(addr_ctx->sock_pair[1]);
addr_ctx->sock_pair[1] = CURL_SOCKET_BAD;
}
#endif
/* Remove socket from event monitoring */
if(addr_ctx->sock_pair[0] != CURL_SOCKET_BAD) {
Curl_multi_will_close(data, addr_ctx->sock_pair[0]);
wakeup_close(addr_ctx->sock_pair[0]);
addr_ctx->sock_pair[0] = CURL_SOCKET_BAD;
}
}
}
#endif
destroy = !addr_ctx->ref_count;
Curl_mutex_release(&addr_ctx->mutx);
if(destroy) {
@ -176,6 +141,12 @@ static void addr_ctx_unlink(struct async_thrdd_addr_ctx **paddr_ctx,
free(addr_ctx->hostname);
if(addr_ctx->res)
Curl_freeaddrinfo(addr_ctx->res);
#ifndef CURL_DISABLE_SOCKETPAIR
#ifndef USE_EVENTFD
wakeup_close(addr_ctx->sock_pair[1]);
#endif
wakeup_close(addr_ctx->sock_pair[0]);
#endif
free(addr_ctx);
}
*paddr_ctx = NULL;
@ -193,10 +164,6 @@ addr_ctx_create(struct Curl_easy *data,
addr_ctx->thread_hnd = curl_thread_t_null;
addr_ctx->port = port;
#ifndef CURL_DISABLE_SOCKETPAIR
addr_ctx->sock_pair[0] = CURL_SOCKET_BAD;
addr_ctx->sock_pair[1] = CURL_SOCKET_BAD;
#endif
addr_ctx->ref_count = 1;
#ifdef HAVE_GETADDRINFO
@ -260,9 +227,7 @@ static CURL_THREAD_RETURN_T CURL_STDCALL getaddrinfo_thread(void *arg)
#pragma clang diagnostic ignored "-Wextra-semi-stmt"
#endif
Curl_thread_cancel_deferred();
Curl_thread_push_cleanup(async_thrd_cleanup, addr_ctx);
Curl_thread_disable_cancel();
Curl_mutex_acquire(&addr_ctx->mutx);
do_abort = addr_ctx->do_abort;
@ -272,7 +237,6 @@ static CURL_THREAD_RETURN_T CURL_STDCALL getaddrinfo_thread(void *arg)
char service[12];
int rc;
Curl_thread_enable_cancel();
#ifdef DEBUGBUILD
Curl_resolve_test_delay();
#endif
@ -289,9 +253,27 @@ static CURL_THREAD_RETURN_T CURL_STDCALL getaddrinfo_thread(void *arg)
else {
Curl_addrinfo_set_port(addr_ctx->res, addr_ctx->port);
}
Curl_mutex_acquire(&addr_ctx->mutx);
do_abort = addr_ctx->do_abort;
Curl_mutex_release(&addr_ctx->mutx);
#ifndef CURL_DISABLE_SOCKETPAIR
if(!do_abort) {
#ifdef USE_EVENTFD
const uint64_t buf[1] = { 1 };
#else
const char buf[1] = { 1 };
#endif
/* Thread is done, notify transfer */
if(wakeup_write(addr_ctx->sock_pair[1], buf, sizeof(buf)) < 0) {
/* update sock_error to errno */
addr_ctx->sock_error = SOCKERRNO;
}
}
#endif
}
Curl_thread_disable_cancel();
Curl_thread_pop_cleanup();
#if defined(__clang__)
#pragma clang diagnostic pop
@ -318,16 +300,13 @@ static CURL_THREAD_RETURN_T CURL_STDCALL gethostbyname_thread(void *arg)
#pragma clang diagnostic ignored "-Wextra-semi-stmt"
#endif
Curl_thread_cancel_deferred();
Curl_thread_push_cleanup(async_thrd_cleanup, addr_ctx);
Curl_thread_disable_cancel();
Curl_mutex_acquire(&addr_ctx->mutx);
do_abort = addr_ctx->do_abort;
Curl_mutex_release(&addr_ctx->mutx);
if(!do_abort) {
Curl_thread_enable_cancel();
if(!do_abort) {
#ifdef DEBUGBUILD
Curl_resolve_test_delay();
#endif
@ -338,9 +317,26 @@ static CURL_THREAD_RETURN_T CURL_STDCALL gethostbyname_thread(void *arg)
if(addr_ctx->sock_error == 0)
addr_ctx->sock_error = RESOLVER_ENOMEM;
}
Curl_mutex_acquire(&addr_ctx->mutx);
do_abort = addr_ctx->do_abort;
Curl_mutex_release(&addr_ctx->mutx);
#ifndef CURL_DISABLE_SOCKETPAIR
if(!do_abort) {
#ifdef USE_EVENTFD
const uint64_t buf[1] = { 1 };
#else
const char buf[1] = { 1 };
#endif
/* Thread is done, notify transfer */
if(wakeup_write(addr_ctx->sock_pair[1], buf, sizeof(buf)) < 0) {
/* update sock_error to errno */
addr_ctx->sock_error = SOCKERRNO;
}
}
#endif
}
Curl_thread_disable_cancel();
Curl_thread_pop_cleanup();
#if defined(__clang__)
#pragma clang diagnostic pop
@ -372,8 +368,14 @@ static void async_thrdd_destroy(struct Curl_easy *data)
bool done;
Curl_mutex_acquire(&addr->mutx);
#ifndef CURL_DISABLE_SOCKETPAIR
if(!addr->do_abort)
Curl_multi_will_close(data, addr->sock_pair[0]);
#endif
addr->do_abort = TRUE;
done = addr->thrd_done;
Curl_mutex_release(&addr->mutx);
if(done) {
Curl_thread_join(&addr->thread_hnd);
CURL_TRC_DNS(data, "async_thrdd_destroy, thread joined");
@ -522,24 +524,16 @@ static void async_thrdd_shutdown(struct Curl_easy *data)
return;
Curl_mutex_acquire(&addr_ctx->mutx);
#ifndef CURL_DISABLE_SOCKETPAIR
if(!addr_ctx->do_abort)
Curl_multi_will_close(data, addr_ctx->sock_pair[0]);
#endif
addr_ctx->do_abort = TRUE;
done = addr_ctx->thrd_done;
#ifndef CURL_DISABLE_SOCKETPAIR
/* We are no longer interested in wakeups */
if(addr_ctx->sock_pair[1] != CURL_SOCKET_BAD) {
wakeup_close(addr_ctx->sock_pair[1]);
addr_ctx->sock_pair[1] = CURL_SOCKET_BAD;
}
#endif
Curl_mutex_release(&addr_ctx->mutx);
DEBUGASSERT(addr_ctx->thread_hnd != curl_thread_t_null);
if(!done && (addr_ctx->thread_hnd != curl_thread_t_null)) {
timediff_t alive_ms = curlx_timediff(curlx_now(), addr_ctx->start);
/* give thread a startup chance to get cancel mode, etc. set up
* before we cancel it. */
if(alive_ms < 5)
curlx_wait_ms(5 - alive_ms);
CURL_TRC_DNS(data, "cancelling resolve thread");
(void)Curl_thread_cancel(&addr_ctx->thread_hnd);
}
@ -721,6 +715,7 @@ CURLcode Curl_async_pollset(struct Curl_easy *data, struct easy_pollset *ps)
{
struct async_thrdd_ctx *thrdd = &data->state.async.thrdd;
CURLcode result = CURLE_OK;
bool thrd_done;
#if !defined(USE_HTTPSRR_ARES) && defined(CURL_DISABLE_SOCKETPAIR)
(void)ps;
@ -736,13 +731,15 @@ CURLcode Curl_async_pollset(struct Curl_easy *data, struct easy_pollset *ps)
if(!thrdd->addr)
return result;
Curl_mutex_acquire(&thrdd->addr->mutx);
thrd_done = thrdd->addr->thrd_done;
Curl_mutex_release(&thrdd->addr->mutx);
if(!thrd_done) {
#ifndef CURL_DISABLE_SOCKETPAIR
/* return read fd to client for polling the DNS resolution status */
if(thrdd->addr->sock_pair[0] != CURL_SOCKET_BAD) {
result = Curl_pollset_add_in(data, ps, thrdd->addr->sock_pair[0]);
}
#else
{
timediff_t milli;
timediff_t ms = curlx_timediff(curlx_now(), thrdd->addr->start);
if(ms < 3)
@ -754,8 +751,8 @@ CURLcode Curl_async_pollset(struct Curl_easy *data, struct easy_pollset *ps)
else
milli = 200;
Curl_expire(data, milli, EXPIRE_ASYNC_NAME);
}
#endif
}
return result;
}

View file

@ -75,14 +75,11 @@ int Curl_thread_cancel(curl_thread_t *hnd);
pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL)
#define Curl_thread_disable_cancel() \
pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL)
#define Curl_thread_cancel_deferred() \
pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL)
#else
#define Curl_thread_push_cleanup(a,b) ((void)a,(void)b)
#define Curl_thread_pop_cleanup() Curl_nop_stmt
#define Curl_thread_enable_cancel() Curl_nop_stmt
#define Curl_thread_disable_cancel() Curl_nop_stmt
#define Curl_thread_cancel_deferred() Curl_nop_stmt
#endif
#endif /* USE_THREADS_POSIX || USE_THREADS_WIN32 */