This is an automated email from the ASF dual-hosted git repository. asf-gitbox-commits pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/qpid-proton.git
commit ec5742577cb08bf8c707bcf5bb212ab4e60d1c79 Author: Andrew Stitcher <[email protected]> AuthorDate: Mon Aug 24 21:44:02 2026 -0400 PROTON-2984: Close connections from the leader on proactor disconnect on_proactor_disconnect() ran pn_connection_driver_close() directly from the uv_walk(), which leader_lead_lh() performs with p->lock dropped. A worker thread can own that connection's driver at the time, so the walk could touch driver state concurrently with the worker. Set pc->disconnecting and let leader_process_pconnection() do the close instead; that only runs for work items the leader owns. Deferring the close exposed a second problem with the disconnect condition. pn_proactor_disconnect() writes p->disconnect_cond under p->lock, but the readers run with the lock dropped: the walk itself, and now the deferred per-connection close as well. Once the leader has cleared p->disconnect a second pn_proactor_disconnect() is no longer suppressed, so it can pn_condition_copy() into disconnect_cond while those readers copy out of it. Both sides go through pn_string, so a reallocation there is a use-after-free rather than merely a torn condition. The listener path had the same unlocked read before this change. Take a copy of the condition under the lock when the disconnect request is consumed, and have both readers use that. Only the leader touches the copy, so they need no lock. The condition race is only reachable when an application calls pn_proactor_disconnect() more than once, so it is latent for single-disconnect users. Assisted-By: Claude Opus 5 <[email protected]> --- c/src/proactor/libuv.c | 30 ++++++++++++++++++++++++------ 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/c/src/proactor/libuv.c b/c/src/proactor/libuv.c index e1acdb33d..59cd1c653 100644 --- a/c/src/proactor/libuv.c +++ b/c/src/proactor/libuv.c @@ -172,6 +172,7 @@ typedef struct pconnection_t { uv_connect_t connect; /* Outgoing connection only */ int connected; /* 0: not connected, <0: connecting after error, 1 = connected ok */ bool connect_started; /* Outgoing connect sequence has been started */ + bool disconnecting; /* Close requested by pn_proactor_disconnect() */ lsocket_t *lsocket; /* Incoming connection only */ @@ -251,6 +252,11 @@ struct pn_proactor_t { size_t active; /* connection/listener count for INACTIVE events */ pn_condition_t *disconnect_cond; /* disconnect condition */ + /* Leader thread only: snapshot of disconnect_cond taken under lock when the + disconnect request is consumed, so the walk and the deferred per-connection + close can read it without racing a later pn_proactor_disconnect(). */ + pn_condition_t *leader_disconnect_cond; + bool has_leader; /* A thread is working as leader */ bool disconnect; /* disconnect requested */ bool batch_working; /* batch is being processed in a worker thread */ @@ -910,6 +916,15 @@ static void check_wake(pconnection_t *pc) { /* Process a pconnection, return true if it has events for a worker thread */ static bool leader_process_pconnection(pconnection_t *pc) { /* Important to do the following steps in order */ + if (pc->disconnecting) { + pc->disconnecting = false; + pn_condition_t *cond = pc->work.proactor->leader_disconnect_cond; + if (cond) { + pn_condition_copy(pn_transport_condition(pc->driver.transport), cond); + } + pn_connection_driver_close(&pc->driver); + } + if (!pc->connected) { /* Retries are driven by the libuv connect callbacks. Restarting here would re-run uv_tcp_init() on the live uv_tcp_t, re-initializing its watcher @@ -991,18 +1006,14 @@ static void on_proactor_disconnect(uv_handle_t* h, void* v) { switch (*(struct_type*)h->data) { case T_CONNECTION: { pconnection_t *pc = (pconnection_t*)h->data; - pn_condition_t *cond = pc->work.proactor->disconnect_cond; - if (cond) { - pn_condition_copy(pn_transport_condition(pc->driver.transport), cond); - } - pn_connection_driver_close(&pc->driver); + pc->disconnecting = true; work_notify(&pc->work); break; } case T_LSOCKET: { pn_listener_t *l = ((lsocket_t*)h->data)->parent; if (l) { - pn_condition_t *cond = l->work.proactor->disconnect_cond; + pn_condition_t *cond = l->work.proactor->leader_disconnect_cond; if (cond) { pn_condition_copy(pn_listener_condition(l), cond); } @@ -1028,6 +1039,11 @@ static pn_event_batch_t *leader_lead_lh(pn_proactor_t *p, uv_run_mode mode) { if (p->disconnect) { p->disconnect = false; if (p->active) { + /* Snapshot while we still hold the lock: from here on only the leader + touches leader_disconnect_cond, so the readers below need no lock. */ + if (p->leader_disconnect_cond && p->disconnect_cond) { + pn_condition_copy(p->leader_disconnect_cond, p->disconnect_cond); + } uv_mutex_unlock(&p->lock); uv_walk(&p->loop, on_proactor_disconnect, NULL); uv_mutex_lock(&p->lock); @@ -1276,6 +1292,7 @@ pn_proactor_t *pn_proactor(void) { uv_timer_init(&p->loop, &p->timer); p->timer.data = p; p->disconnect_cond = pn_condition(); + p->leader_disconnect_cond = pn_condition(); return p; } @@ -1297,6 +1314,7 @@ void pn_proactor_free(pn_proactor_t *p) { uv_cond_destroy(&p->cond); pn_collector_free(p->collector); pn_condition_free(p->disconnect_cond); + pn_condition_free(p->leader_disconnect_cond); free(p); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
