chenBright commented on code in PR #3509:
URL: https://github.com/apache/brpc/pull/3509#discussion_r4226438597


##########
src/brpc/ubshm/ub_ring_manager.cpp:
##########
@@ -68,57 +80,191 @@ RETURN_CODE UBRingManager::UbrMgrInit() {
     g_ubr_mgr.trx_mgr = (UbrTrx *)malloc(trx_mgr_size);
     size_t trx_mgr_status_size = g_ubr_mgr.trx_cap * sizeof(UbrMgrUnitStatus);
     g_ubr_mgr.trx_mgr_unit_status = (UbrMgrUnitStatus 
*)malloc(trx_mgr_status_size);
-    if (UNLIKELY(g_ubr_mgr.trx_mgr == nullptr ||
-                 g_ubr_mgr.trx_mgr_unit_status == nullptr)) {
+    size_t trx_mgr_id_size = g_ubr_mgr.trx_cap * sizeof(uint64_t);
+    g_ubr_mgr.trx_mgr_unit_id = (uint64_t *)malloc(trx_mgr_id_size);
+    size_t trx_mgr_ctl_size = g_ubr_mgr.trx_cap * sizeof(UbrCleanupCtl *);
+    g_ubr_mgr.trx_mgr_unit_ctl = (UbrCleanupCtl **)malloc(trx_mgr_ctl_size);
+    if (BAIDU_UNLIKELY(g_ubr_mgr.trx_mgr == nullptr ||
+                 g_ubr_mgr.trx_mgr_unit_status == nullptr ||
+                 g_ubr_mgr.trx_mgr_unit_id == nullptr ||
+                 g_ubr_mgr.trx_mgr_unit_ctl == nullptr)) {
         LOG(ERROR) << "Ubr manager memory allocation failed.";
         UbrMgrFini();
         return UBRING_ERR;
     }
 
-    memset(g_ubr_mgr.trx_mgr, 0, trx_mgr_size);
+    // UbrTrx holds butil::atomic members, so it is not trivially copyable and
+    // must not be memset: value-initialize every slot instead. A slot is
+    // re-initialized by AcquireUbrTrxFromMgr before it is handed out, and no
+    // code reads a slot that was never acquired.
+    for (uint32_t i = 0; i < g_ubr_mgr.trx_cap; ++i) {
+        new (&g_ubr_mgr.trx_mgr[i]) UbrTrx();
+    }
     memset(g_ubr_mgr.trx_mgr_unit_status, UBR_MGR_UNIT_FREE, 
trx_mgr_status_size);
+    memset(g_ubr_mgr.trx_mgr_unit_id, 0, trx_mgr_id_size);
+    memset(g_ubr_mgr.trx_mgr_unit_ctl, 0, trx_mgr_ctl_size);
     LinkInfoInit();
     return UBRING_OK;
 }
 
 void UBRingManager::UbrMgrFini() {
+    // Refuse new pool users and new acquisitions first, unconditionally: the
+    // partial-init path below (a failed allocation) must not let a callback
+    // walk arrays that the failure left null either.
+    {
+        BAIDU_SCOPED_LOCK(g_ubr_trx_mgr_mtx);
+        g_ubr_mgr.shutting_down = true;
+    }
+
+    // The pool arrays are only safe to walk once UbrMgrInit has allocated and
+    // zeroed all four of them; its allocation-failure path calls this function
+    // with some of them still null and the others uninitialized.
+    const bool pool_ready = g_ubr_mgr.trx_mgr != nullptr &&
+                            g_ubr_mgr.trx_mgr_unit_status != nullptr &&
+                            g_ubr_mgr.trx_mgr_unit_id != nullptr &&
+                            g_ubr_mgr.trx_mgr_unit_ctl != nullptr;
+
+    // Wait out close/heartbeat callbacks that are already dispatched, then
+    // stop the timers, before the pool is freed: UbrTimerDel alone is
+    // non-blocking, so a callback that already passed its generation check
+    // could still be reading UbrTrx when FREE_PTR(trx_mgr) runs.
+    // UbrTimerDelAndWait blocks until a dispatched callback has returned, but
+    // it must run without g_ubr_trx_mgr_mtx because the callbacks take that
+    // lock themselves; snapshot the slots under the lock first.
+    if (pool_ready) {
+        std::vector<uint32_t> used_slots;
+        {
+            BAIDU_SCOPED_LOCK(g_ubr_trx_mgr_mtx);
+            for (uint32_t i = 0; i < g_ubr_mgr.trx_cap; ++i) {
+                if (g_ubr_mgr.trx_mgr_unit_status[i] == UBR_MGR_UNIT_USED) {
+                    used_slots.push_back(i);
+                }
+            }
+        }
+        // Wait for the callback-driven pool accesses that are already running
+        // (they read UbrTrx without owning a timer), then for the timers. Both
+        // waits must happen without g_ubr_trx_mgr_mtx: the close/heartbeat
+        // callbacks and the accesses themselves take that lock.
+        for (;;) {
+            {
+                BAIDU_SCOPED_LOCK(g_ubr_trx_mgr_mtx);
+                if (g_ubr_mgr.active_pool_ops == 0) {
+                    break;
+                }
+            }
+            usleep(1000);
+        }
+        // No timer can be armed after the flag: arming and taking the flag are
+        // serialized by the manager lock (ArmTimersExclusive), and a cleanup
+        // control object can no longer be published either
+        // (TryPublishUnitCleanupCtl refuses once shutting down). So the timers
+        // armed before the flag are the complete set, and stopping them all
+        // leaves the pool quiescent.
+        for (uint32_t i : used_slots) {
+            UbrTimerDelAndWait(&g_ubr_mgr.trx_mgr[i].close_timer);
+            UbrTimerDelAndWait(&g_ubr_mgr.trx_mgr[i].hb_timer);
+        }
+    }
+
+    // Cancel the pending delayed cleanups and wait for the in-flight ones
+    // (each holds one extra reference) to finish, before the pool memory
+    // they touch is freed. A ctl whose timer is still starting can only be
+    // cancelled in a later round, hence the retry-to-stability loop.
+    bool busy = true;
+    while (busy) {
+        busy = false;
+        {
+            BAIDU_SCOPED_LOCK(g_ubr_trx_mgr_mtx);
+            if (pool_ready) {
+                for (uint32_t i = 0; i < g_ubr_mgr.trx_cap; ++i) {
+                    UbrCleanupCtl* ctl = g_ubr_mgr.trx_mgr_unit_ctl[i];
+                    if (ctl == nullptr) {
+                        continue;
+                    }
+                    if (UbrTimerDel(&ctl->timer) == 0) {
+                        ctl->ReleaseRef();   // timer/callback reference

Review Comment:
   There is still a lifetime race when a one-shot timer has been dispatched but 
`UbrTimerOnFire` has not yet performed its handle-slot CAS:
   
   1. `UbrTimerDel` takes `ctl->timer`.
   2. The underlying `bthread_timer_del` returns `1`, but `UbrTimerDel` returns 
`0` because it won the slot.
   3. This path releases the timer/callback reference. Once only the manager 
anchor remains, `UbrMgrFini` can release that anchor and delete `ctl`.
   4. The dispatched wrapper resumes and performs its CAS through `task->slot`, 
which points to the freed `ctl->timer`.
   
   Winning the slot suppresses the user callback, but the wrapper still 
accesses the slot before returning. The remaining schedule reference protects 
`UbrTimerTask`; it does not keep `UbrCleanupCtl` alive. Force-close followed by 
slot reuse can expose the same race when the old manager anchor is retired.



##########
src/brpc/ubshm/shm/shm_ubs.cpp:
##########
@@ -410,58 +426,78 @@ static void DeleteShmToList(ShmList* shm_list)
     shm_list->size--;
 }
 
-void *UbsShmCallback(void* args)
+void *UbsShmCallback(void* args, uint64_t)
 {
     ShmList *shm_list = (ShmList*)args;
-    if (UNLIKELY(shm_list == nullptr)) {
+    if (BAIDU_UNLIKELY(shm_list == nullptr)) {
         LOG(ERROR) << "Shm list is null.";
         return nullptr;
     }
 
-    LOCK_GUARD(shm_list->shm_lock);
-    while (shm_list->head != nullptr) {
-        SHM shm = shm_list->head->shm;
-        if (shm.addr == nullptr) {
-            LOG(ERROR) << "Ubs input shm param is invalid, addr is NULL.";
+    // Drain one node per fire and keep the SDK calls outside the lock, so
+    // a slow daemon cannot stall the timer thread for the whole list.
+    SHM shm;
+    {
+        BAIDU_SCOPED_LOCK(shm_list->shm_lock);
+        if (shm_list->head == nullptr) {
             return nullptr;
         }
+        shm = shm_list->head->shm;
+    }
+    if (shm.addr == nullptr) {
+        LOG(ERROR) << "Ubs input shm param is invalid, addr is NULL.";
+        BAIDU_SCOPED_LOCK(shm_list->shm_lock);
+        DeleteShmToList(shm_list);
+        return nullptr;
+    }
 
-        int ret = ubsmem_shmem_unmap(shm.addr, shm.len);
-        if (ret != UBSM_OK) {
-            if (ret == UBSM_ERR_NET) {
-                return nullptr;
-            }
-            LOG(ERROR) << "Ubs unmap shm=" << shm.name << " length=" << 
shm.len << " failed, ret=" << ret;
-            return nullptr;
+    int ret = ubsmem_shmem_unmap(shm.addr, shm.len);
+    if (ret != UBSM_OK) {
+        if (ret == UBSM_ERR_NET) {
+            return nullptr;              // retried on the next fire
         }
-        LOG(INFO) << "Ubs unmap shm=" << shm.name << " length=" << shm.len << 
" success.";
+        LOG(ERROR) << "Ubs unmap shm=" << shm.name << " length=" << shm.len << 
" failed, ret=" << ret;
+        return nullptr;                  // node stays at head, retried
+    }
 
-        ret = ubsmem_shmem_deallocate(shm.name);
-        if (ret != UBSM_OK) {
-            DeleteShmToList(shm_list);
-            LOG(ERROR) << "Ubs delete shm=" << shm.name << " failed, ret=" << 
ret;
-            return nullptr;
-        }
+    {
+        BAIDU_SCOPED_LOCK(shm_list->shm_lock);
         DeleteShmToList(shm_list);
-        LOG(INFO) << "Ubs free local shm=" << shm.name << " length=" << 
shm.len << " success.";
     }
 
+    ret = ubsmem_shmem_deallocate(shm.name);
+    if (ret != UBSM_OK) {
+        LOG(ERROR) << "Ubs delete shm=" << shm.name << " failed, ret=" << ret;
+    }
     return nullptr;
 }
 
 RETURN_CODE UbsShmAddTimer(ShmList *shm_list)
 {
-    const uint32_t timer_interval_s = FLAGS_ub_flying_io_timeout_s;
-    itimerspec time_spec = {
-        .it_interval = {.tv_sec = timer_interval_s, .tv_nsec = 0},
-        .it_value = {.tv_sec = 0, .tv_nsec = 1}
-    };
-    int timer_fd = TimerStart(&time_spec, UbsShmCallback, (void*)shm_list);
-    if (UNLIKELY(timer_fd == -1)) {
+    // This timer is periodic -- it drains one pending unmap per fire -- while
+    // the same flag is also the one-shot delay of the delayed-clear path, 
where
+    // 0 legitimately means "do not wait for in-flight IO". A periodic timer
+    // needs a positive period, so fall back to a drain period of one second
+    // instead of either refusing to start (which would break UBS shm init for 
a
+    // valid tuning of the other use) or degrading to a single drain.
+    uint64_t period_s = (uint64_t)FLAGS_ub_flying_io_timeout_s;
+    if (BAIDU_UNLIKELY(FLAGS_ub_flying_io_timeout_s <= 0)) {
+        LOG(WARNING) << "ub_flying_io_timeout_s=" << 
FLAGS_ub_flying_io_timeout_s
+                     << " disables the delayed-clear wait; using a 1s drain 
period "
+                     << "for the shm cleanup timer.";
+        period_s = 1;
+    }
+    const uint64_t timer_interval_us = period_s * SEC_TO_USEC;
+    // The generation argument is the caller's staleness guard; UbsShmCallback
+    // ignores it and never frees `shm_list', whose lifetime is settled by
+    // DestroyShmTimer waiting through UbrTimerDelAndWait before it tears the
+    // list down. 0 is therefore deliberate here, not a missing generation.
+    RETURN_CODE rc = UbrTimerStartPeriodic(&g_shm_timer_id, 0, 
timer_interval_us,
+                                           UbsShmCallback, (void*)shm_list, 0);
+    if (BAIDU_UNLIKELY(rc != UBRING_OK)) {

Review Comment:
   This schedules `UbsShmCallback` on the process-wide bthread timer thread, 
where it still directly calls `ubsmem_shmem_unmap` and 
`ubsmem_shmem_deallocate`.
   
   Moving these calls outside `shm_lock` and processing one node per invocation 
reduces lock contention and batch size. A single slow SDK call can still occupy 
the global timer thread for its entire duration.
   
   Previously, this work ran on a dedicated UBRing timer thread. With this 
change, a slow daemon or network failure can also delay ordinary RPC timeouts, 
backup requests, and bthread sleep/timed-wait wakeups across the process. The 
delayed transaction cleanup callbacks perform synchronous SDK cleanup on the 
same thread as well.



##########
src/brpc/ubshm/ub_ring.cpp:
##########
@@ -89,25 +235,53 @@ RETURN_CODE UBRing::UbrTrxClose() {
             LOG(WARNING) << "Local shm " << _trx->local_shm.name
             << " wait for the peer to close timed out, force cleanup.";
             _trx->ubr_rx.trx_state = UBR_STATE_CLOSED;
-            // Force synchronous cleanup instead of relying on async timer
-            DeleteTimerSafe((uint32_t)_trx->timer_fd);
-            DeleteTimerSafe((uint32_t)_trx->hb_timer_fd);
-            if (_trx->ubr_tx.remote_rx_event_q.addr != nullptr) {
-                ((UbrEventQMsg *)_trx->ubr_tx.remote_rx_event_q.addr)->flag = 
UBR_STATE_CLOSED;
+            // Wait out the close/heartbeat callbacks, which may schedule a
+            // delayed cleanup, then settle the cleanup ownership: force
+            // runs the cleanup itself when it can claim it, and leaves it
+            // to an already running delayed-clear callback otherwise.
+            UbrTimerDelAndWait(&_trx->close_timer);
+            UbrTimerDelAndWait(&_trx->hb_timer);
+            UbrCleanupCtl* ctl = 
UBRingManager::SnapshotUnitCleanupCtl(_trx->trx_mgr_index);
+            if (ctl != nullptr && ctl->ubr_id != expect_ubr_id) {
+                ctl->ReleaseRef();               // snapshot reference
+                ctl = nullptr;                   // slot reused, not ours
             }
-            if (UNLIKELY(UbrTrxFreeShm(_trx) != UBRING_OK)) {
-                LOG(WARNING) << "Force close, local shm " << 
_trx->local_shm.name << " free failed.";
+            bool cleanup_owned = false;
+            if (ctl != nullptr) {
+                int expected = UBR_CLEANUP_PENDING;
+                if (ATOMIC_COMPARE_EXCHANGE_STRONG(ctl->state, expected, 
UBR_CLEANUP_RUNNING)) {
+                    cleanup_owned = true;
+                    if (UbrTimerDel(&ctl->timer) == 0) {
+                        ctl->ReleaseRef();   // timer/callback reference
+                    }
+                }
+            } else if (ATOMIC_LOAD(_trx->ubr_id) == expect_ubr_id) {
+                cleanup_owned = true;
             }

Review Comment:
   `SnapshotUnitCleanupCtl` releases the manager lock before this branch claims 
cleanup ownership. A concurrent SDK fault callback can publish a cleanup 
control object after the snapshot returned `nullptr`.
   
   This is reachable while an active close is in progress: `UbrPassiveClearTrx` 
receives `UBRING_REENTRY` and proceeds through `ClearTrxResource`.
   
   The following interleaving allows two cleanup owners:
   
   1. Force-close stops the periodic timers and snapshots a null control object.
   2. The SDK fault callback publishes a new control object for the same 
transaction generation and arms its cleanup timer.
   3. Force-close sets `cleanup_owned = true` using its earlier null snapshot.
   4. The new timer callback independently claims its control object and starts 
cleanup while the transaction slot is still marked USED.
   
   With the explicitly supported `ub_flying_io_timeout_s=0`, or a sufficiently 
slow synchronous cleanup, both paths can unmap the same shared memory or access 
event queues after they have been unmapped.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to