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


##########
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:
   UBS cleanup now runs on a worker, so the normal path no longer calls the SDK 
from the global timer thread. If the worker fails to start, it still falls back 
to inline cleanup.



##########
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:
   The snapshot and forced-cleanup claim now happen under the same manager 
lock. Once force-close claims cleanup, a later attempt to publish delayed 
cleanup is rejected.



##########
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:
   `UbrTimerDel` now waits for an already-dispatched one-shot timer wrapper to 
finish accessing the slot before returning, so the control object can be 
released safely.



-- 
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