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]

Reply via email to