Copilot commented on code in PR #3509:
URL: https://github.com/apache/brpc/pull/3509#discussion_r4131437931
##########
src/brpc/ubshm/timer/timer_mgr.cpp:
##########
@@ -15,454 +15,258 @@
// specific language governing permissions and limitations
// under the License.
-#define _GNU_SOURCE
-#include <pthread.h>
-#include <sched.h>
-#include <errno.h>
-#include <stdio.h>
-#include <stdlib.h>
-#include <unistd.h>
-#include <atomic>
-#include <sys/resource.h>
+#include <new>
+#include "bthread/bthread.h" // bthread_usleep
+#include "bthread/unstable.h" // bthread_timer_add/del
+#include "butil/atomicops.h"
+#include "butil/time.h"
#include "brpc/ubshm/timer/timer_mgr.h"
namespace brpc {
namespace ubring {
-int32_t g_epoll_fd = -1;
-std::atomic<uint32_t> g_total_timer_num(0);
-TimerFdCtx *g_timer_fd_ctx_map = nullptr;
-uint32_t g_max_system_fd = 0;
-static pthread_t g_epoll_execute_thread = 0;
-static int32_t g_timer_module_initialized = 0;
-
-#if defined(OS_MACOSX)
-static int timerfd_create_macosx(int clockid, int flags);
-static int timerfd_settime_macosx(int fd, int flags,
- const itimerspec *new_value,
- itimerspec *old_value);
-#endif
-
-static RETURN_CODE DeleteTimerInner(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
- return UBRING_OK;
- }
-
- if (pthread_spin_lock(&g_timer_fd_ctx_map[fd].spin_lock) != 0) {
- return UBRING_ERR;
- }
-
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
- return UBRING_OK;
- }
-
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].periodical = 0;
- g_timer_fd_ctx_map[fd].fd = 0;
-
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
-
-#if defined(OS_LINUX)
- epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, (int)fd, nullptr);
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr);
-#endif
-
- uint64_t exp = 0;
- read((int)fd, &exp, sizeof(exp));
-
- close((int)fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
- return UBRING_OK;
-}
-
-static RETURN_CODE StartTimeEpoll(void) {
-#if defined(OS_LINUX)
- g_epoll_fd = epoll_create1(0);
-#elif defined(OS_MACOSX)
- g_epoll_fd = kqueue();
-#endif
- if (UNLIKELY(g_epoll_fd == -1)) {
- LOG(ERROR) << "Failed to create epoll/kqueue. errno=" << errno;
- return UBRING_ERR;
- }
-
- int ret = pthread_create(&g_epoll_execute_thread, nullptr, TimerEpoll,
nullptr);
- if (UNLIKELY(ret != 0)) {
- LOG(ERROR) << "Failed to create thread err=" << ret;
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerSpinLocksInit(void) {
- if (g_timer_fd_ctx_map == nullptr) {
- LOG(ERROR) << "Timer module is not fully initialized.";
- return UBRING_ERR;
- }
-
- for (uint32_t fd = 0; fd < g_max_system_fd; fd++) {
- int ret = pthread_spin_init(&g_timer_fd_ctx_map[fd].spin_lock,
- PTHREAD_PROCESS_PRIVATE);
- if (ret != EOK) {
- LOG(ERROR) << "Failed to initialize spin lock for fd=" << fd;
- for (uint32_t cleanup_fd = 0; cleanup_fd < fd; cleanup_fd++) {
-
pthread_spin_destroy(&g_timer_fd_ctx_map[cleanup_fd].spin_lock);
- }
- return UBRING_ERR;
- }
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE ExecuteCallback(int32_t timer_fd) {
- UnifiedCallback((void *)(&g_timer_fd_ctx_map[timer_fd]));
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerCtxMapCompletion(void) {
- memset(g_timer_fd_ctx_map, 0, sizeof(TimerFdCtx) * g_max_system_fd);
-
- RETURN_CODE ret = TimerSpinLocksInit();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init spin locks for timer module.";
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-RETURN_CODE TimerInit(void) {
- if (g_timer_module_initialized > 0) {
- return UBRING_OK;
- }
-
- g_total_timer_num.store(0);
-
- struct rlimit rlim;
- if (getrlimit(RLIMIT_NOFILE, &rlim) != UBRING_OK) {
- LOG(ERROR) << "Failed to get fd";
- return UBRING_ERR;
- }
- g_max_system_fd = (uint32_t)rlim.rlim_cur;
-
- if (g_timer_fd_ctx_map == nullptr) {
- g_timer_fd_ctx_map = (TimerFdCtx *)malloc(sizeof(TimerFdCtx) *
g_max_system_fd);
- if (UNLIKELY(!g_timer_fd_ctx_map)) {
- LOG(ERROR) << "Fail to malloc space for timer modules. errno=%d",
errno;
- return UBRING_ERR;
- }
-
- RETURN_CODE ret = TimerCtxMapCompletion();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init main data structure of Time Module.
ret=" << ret;
- free(g_timer_fd_ctx_map);
- g_timer_fd_ctx_map = nullptr;
- return UBRING_ERR;
- }
- }
-
- RETURN_CODE ret = StartTimeEpoll();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to start Timer Epoll. ret=" << ret;
- if (LIKELY(g_timer_fd_ctx_map != nullptr)) {
- FREE_PTR(g_timer_fd_ctx_map);
+namespace {
+
+enum UbrTimerState {
+ kStarting = 0, // published, not scheduled
yet
+ kScheduled = 1,
+ kDead = 2 // scheduling failed
+};
+
+} // namespace
+
+// Reference rules: one "owner" ref for the handle slot, one "schedule" ref
+// per pending/running bthread schedule, plus one ref held by the starter
+// until its post-schedule bookkeeping is done. The schedule ref is
+// consumed by the firing callback or by the deleter whose
+// bthread_timer_del returned 0 (cancelled before run); the owner ref is
+// consumed by whoever takes the task out of *slot -- a deleter, or the
+// one-shot firing callback itself, which exits the slot BEFORE running
+// the callback so that the callback may free the object storing the slot.
+// All atomics are seq_cst so no interleaving can release a ref twice or
+// free the task while a callback or the starter still touches it.
+struct UbrTimerTask {
+ butil::atomic<UbrTimerId>* slot;
+ butil::atomic<bthread_timer_t> id;
Review Comment:
`UbrTimerId` in the public header is a pointer to
`brpc::ubring::UbrTimerTask`, but this file defines a different `UbrTimerTask`
inside an anonymous namespace. The task pointers used by `slot` therefore have
unrelated types, so the assignments/CAS calls (for example at lines 165-166 and
`TakeOutTask`) do not compile. Define the task in the enclosing `brpc::ubring`
namespace (and keep only helper functions anonymous), or make the opaque type
declaration match the definition.
This issue also appears on line 102 of the same file.
##########
src/brpc/ubshm/ub_ring.cpp:
##########
@@ -47,27 +56,164 @@ UBRing::~UBRing()
RETURN_CODE UBRing::UbrTrxMapShm(SHM *local_shm, SHM *remote_shm)
{
RETURN_CODE rc = UbrTrxMapLocalShm(local_shm);
- if (UNLIKELY(rc != UBRING_OK)) {
+ if (BAIDU_UNLIKELY(rc != UBRING_OK)) {
LOG(ERROR) << "Trx map local shared memory failed.";
return rc;
}
rc = UbrTrxMapRemoteShm(remote_shm);
- if (UNLIKELY(rc != UBRING_OK)) {
+ if (BAIDU_UNLIKELY(rc != UBRING_OK)) {
LOG(ERROR) << "Trx map remote shared memory failed.";
return rc;
}
return UBRING_OK;
}
+// Stop a per-trx timer before its slot becomes reusable. Outside a per-trx
+// timer callback this waits -- UbrTimerDelAndWait returns only after a
+// callback that was already dispatched has left the trx -- which is what lets
+// the caller clear the trx and free its shared memory afterwards. Inside the
+// callback itself the wait would join the task that is currently running, so
+// the non-blocking delete is used instead; no wait is needed there, because
+// bthread dispatches every timer callback from one global timer thread
+// (TimerThread in bthread/timer_thread.{h,cpp} is created by a single
+// pthread_once), so this trx's sibling timer callback cannot be running
+// concurrently and the non-blocking delete already removed it from the timer
+// heap.
+static void UbrStopTrxTimer(butil::atomic<UbrTimerId>* slot,
+ bool in_timer_callback) {
+ if (in_timer_callback) {
+ UbrTimerDel(slot);
+ } else {
+ UbrTimerDelAndWait(slot);
+ }
+}
+
+static void UbrDoAsynClearWork(UbrTrx *trx, uint64_t expect_ubr_id) {
+ if (BAIDU_UNLIKELY(UBRing::UbrTrxFreeShm(trx) != UBRING_OK)) {
+ LOG(ERROR) << "Trx close, wait for local shm " << trx->local_shm.name
<< " free fail.";
+ }
+ if (BAIDU_UNLIKELY(UBRingManager::ReleaseUbrTrxFromMgr(trx, expect_ubr_id)
!= UBRING_OK)) {
+ LOG(ERROR) << "Trx close, release shm " << trx->local_shm.name << "
trx failed.";
+ }
+}
+
+static void UbrDoPassiveClearWork(UbrTrx *trx, uint64_t expect_ubr_id) {
+ int rc = ShmLocalFree(&trx->remote_shm);
+ if (rc != UBRING_OK) {
+ LOG(ERROR) << "Trx passive clear, delete remote shm " <<
trx->remote_shm.name
+ << " failed. ret=" << rc;
+ }
+ rc = ShmLocalFree(&trx->local_shm);
+ if (rc != UBRING_OK) {
+ LOG(ERROR) << "Trx passive clear, delete local shm " <<
trx->local_shm.name
+ << " failed. ret=" << rc;
+ }
+ if (BAIDU_UNLIKELY(UBRingManager::ReleaseUbrTrxFromMgr(trx, expect_ubr_id)
!= UBRING_OK)) {
+ LOG(ERROR) << "Trx passive clear, release shm " << trx->local_shm.name
<< " trx failed.";
+ }
+}
+
+// Schedule the delayed cleanup of `trx'. The cleanup ownership lives in the
+// per-acquisition control object, so exactly one of the delayed-clear
+// callback and a force close ever runs the cleanup. `work' is the cleanup
+// body, used directly when the timer cannot be started.
+static RETURN_CODE UbrScheduleClearTimer(UbrTrx *trx, uint64_t expect_ubr_id,
+ void* (*cb)(void*, uint64_t),
+ void (*work)(UbrTrx*, uint64_t)) {
+ if (BAIDU_UNLIKELY(trx == nullptr || trx->local_shm.addr == nullptr)) {
+ return UBRING_OK; // released trx, stale event
+ }
+ // A callback that outlived its generation must not capture the id of the
+ // slot's new occupant, nor schedule cleanup for that new transaction.
+ if (BAIDU_UNLIKELY(ATOMIC_LOAD(trx->ubr_id) != expect_ubr_id)) {
+ return UBRING_OK; // stale event on a reused slot
+ }
+ if (trx->cleanup_ctl.load() != nullptr) {
+ return UBRING_OK; // cleanup already scheduled
+ }
+ auto* ctl = new (std::nothrow) UbrCleanupCtl();
+ if (BAIDU_UNLIKELY(ctl == nullptr)) {
+ LOG(ERROR) << "Fail to malloc ubr cleanup ctl.";
+ return UBRING_ERR;
+ }
+ ctl->trx = trx;
+ ctl->ubr_id = expect_ubr_id;
+ ctl->state.store(UBR_CLEANUP_PENDING);
+ ctl->timer = nullptr;
+ ctl->ref.store(2); // timer/callback + starter; the
+ // manager anchor is taken by
+ // TryPublishUnitCleanupCtl
+
+ UbrCleanupCtl* expected = nullptr;
+ if (!trx->cleanup_ctl.compare_exchange_strong(expected, ctl)) {
+ delete ctl; // another schedule won
+ return UBRING_OK;
+ }
+ if (!UBRingManager::TryPublishUnitCleanupCtl(trx->trx_mgr_index,
+ ctl->ubr_id, ctl)) {
+ // The slot was released (and possibly reused) before we could
+ // anchor: force close or the new occupant owns it now. Nothing
+ // was armed yet -- just undo the trx-side publication.
+ expected = ctl;
+ trx->cleanup_ctl.compare_exchange_strong(expected, nullptr);
Review Comment:
Publishing `ctl` into `trx->cleanup_ctl` and anchoring it in the manager are
separate operations. A force-close can observe no manager anchor between these
lines, claim the no-ctl cleanup path, release/reuse the slot, and then this
thread performs the CAS through the reused `UbrTrx`; the generation check in
`TryPublishUnitCleanupCtl` happens too late to make that field access safe.
Publish the control object and validate the generation atomically under the
manager lock, or keep the publication solely in the manager-owned array.
This issue also appears in the following locations of the same file:
- line 190
- line 453
- line 487
##########
src/brpc/ubshm/ub_ring_manager.cpp:
##########
@@ -68,57 +78,154 @@ 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);
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() {
+ // 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);
+ }
+ }
+ }
+ for (uint32_t i : used_slots) {
+ UbrTimerDelAndWait(&g_ubr_mgr.trx_mgr[i].close_timer);
+ UbrTimerDelAndWait(&g_ubr_mgr.trx_mgr[i].hb_timer);
Review Comment:
This shutdown loop has the same slot-visibility hole: a running
close/heartbeat callback can remove its periodic task from the slot via
`UbrTimerDel`, after which `UbrTimerDelAndWait` returns without waiting.
`UbrMgrFini` can then free `trx_mgr` while that callback still uses the pooled
`UbrTrx`. Shutdown needs an independent in-flight callback count/reference, not
just the timer slot, before freeing the pool.
##########
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);
Review Comment:
`UbrTimerDelAndWait` cannot establish the quiescence promised by this
teardown path when the callback has already entered `UbrTrxCloseCallback`: that
callback calls `UbrStopTrxTimer(..., true)`, which removes its own periodic
task from the slot before it continues into `ClearTrxResource`. A concurrent
force close therefore sees both slots as null and returns from these calls
immediately, then can free/release the trx while the callback is still
accessing it and scheduling delayed cleanup. Track in-flight callbacks
independently of the slot (or defer self-removal until callback exit) before
allowing the force-close path to reclaim the trx.
##########
src/brpc/ubshm/timer/timer_mgr.cpp:
##########
@@ -15,454 +15,258 @@
// specific language governing permissions and limitations
// under the License.
-#define _GNU_SOURCE
-#include <pthread.h>
-#include <sched.h>
-#include <errno.h>
-#include <stdio.h>
-#include <stdlib.h>
-#include <unistd.h>
-#include <atomic>
-#include <sys/resource.h>
+#include <new>
+#include "bthread/bthread.h" // bthread_usleep
+#include "bthread/unstable.h" // bthread_timer_add/del
+#include "butil/atomicops.h"
+#include "butil/time.h"
#include "brpc/ubshm/timer/timer_mgr.h"
namespace brpc {
namespace ubring {
-int32_t g_epoll_fd = -1;
-std::atomic<uint32_t> g_total_timer_num(0);
-TimerFdCtx *g_timer_fd_ctx_map = nullptr;
-uint32_t g_max_system_fd = 0;
-static pthread_t g_epoll_execute_thread = 0;
-static int32_t g_timer_module_initialized = 0;
-
-#if defined(OS_MACOSX)
-static int timerfd_create_macosx(int clockid, int flags);
-static int timerfd_settime_macosx(int fd, int flags,
- const itimerspec *new_value,
- itimerspec *old_value);
-#endif
-
-static RETURN_CODE DeleteTimerInner(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
- return UBRING_OK;
- }
-
- if (pthread_spin_lock(&g_timer_fd_ctx_map[fd].spin_lock) != 0) {
- return UBRING_ERR;
- }
-
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
- return UBRING_OK;
- }
-
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].periodical = 0;
- g_timer_fd_ctx_map[fd].fd = 0;
-
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
-
-#if defined(OS_LINUX)
- epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, (int)fd, nullptr);
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr);
-#endif
-
- uint64_t exp = 0;
- read((int)fd, &exp, sizeof(exp));
-
- close((int)fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
- return UBRING_OK;
-}
-
-static RETURN_CODE StartTimeEpoll(void) {
-#if defined(OS_LINUX)
- g_epoll_fd = epoll_create1(0);
-#elif defined(OS_MACOSX)
- g_epoll_fd = kqueue();
-#endif
- if (UNLIKELY(g_epoll_fd == -1)) {
- LOG(ERROR) << "Failed to create epoll/kqueue. errno=" << errno;
- return UBRING_ERR;
- }
-
- int ret = pthread_create(&g_epoll_execute_thread, nullptr, TimerEpoll,
nullptr);
- if (UNLIKELY(ret != 0)) {
- LOG(ERROR) << "Failed to create thread err=" << ret;
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerSpinLocksInit(void) {
- if (g_timer_fd_ctx_map == nullptr) {
- LOG(ERROR) << "Timer module is not fully initialized.";
- return UBRING_ERR;
- }
-
- for (uint32_t fd = 0; fd < g_max_system_fd; fd++) {
- int ret = pthread_spin_init(&g_timer_fd_ctx_map[fd].spin_lock,
- PTHREAD_PROCESS_PRIVATE);
- if (ret != EOK) {
- LOG(ERROR) << "Failed to initialize spin lock for fd=" << fd;
- for (uint32_t cleanup_fd = 0; cleanup_fd < fd; cleanup_fd++) {
-
pthread_spin_destroy(&g_timer_fd_ctx_map[cleanup_fd].spin_lock);
- }
- return UBRING_ERR;
- }
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE ExecuteCallback(int32_t timer_fd) {
- UnifiedCallback((void *)(&g_timer_fd_ctx_map[timer_fd]));
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerCtxMapCompletion(void) {
- memset(g_timer_fd_ctx_map, 0, sizeof(TimerFdCtx) * g_max_system_fd);
-
- RETURN_CODE ret = TimerSpinLocksInit();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init spin locks for timer module.";
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-RETURN_CODE TimerInit(void) {
- if (g_timer_module_initialized > 0) {
- return UBRING_OK;
- }
-
- g_total_timer_num.store(0);
-
- struct rlimit rlim;
- if (getrlimit(RLIMIT_NOFILE, &rlim) != UBRING_OK) {
- LOG(ERROR) << "Failed to get fd";
- return UBRING_ERR;
- }
- g_max_system_fd = (uint32_t)rlim.rlim_cur;
-
- if (g_timer_fd_ctx_map == nullptr) {
- g_timer_fd_ctx_map = (TimerFdCtx *)malloc(sizeof(TimerFdCtx) *
g_max_system_fd);
- if (UNLIKELY(!g_timer_fd_ctx_map)) {
- LOG(ERROR) << "Fail to malloc space for timer modules. errno=%d",
errno;
- return UBRING_ERR;
- }
-
- RETURN_CODE ret = TimerCtxMapCompletion();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init main data structure of Time Module.
ret=" << ret;
- free(g_timer_fd_ctx_map);
- g_timer_fd_ctx_map = nullptr;
- return UBRING_ERR;
- }
- }
-
- RETURN_CODE ret = StartTimeEpoll();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to start Timer Epoll. ret=" << ret;
- if (LIKELY(g_timer_fd_ctx_map != nullptr)) {
- FREE_PTR(g_timer_fd_ctx_map);
+namespace {
+
+enum UbrTimerState {
+ kStarting = 0, // published, not scheduled
yet
+ kScheduled = 1,
+ kDead = 2 // scheduling failed
+};
+
+} // namespace
+
+// Reference rules: one "owner" ref for the handle slot, one "schedule" ref
+// per pending/running bthread schedule, plus one ref held by the starter
+// until its post-schedule bookkeeping is done. The schedule ref is
+// consumed by the firing callback or by the deleter whose
+// bthread_timer_del returned 0 (cancelled before run); the owner ref is
+// consumed by whoever takes the task out of *slot -- a deleter, or the
+// one-shot firing callback itself, which exits the slot BEFORE running
+// the callback so that the callback may free the object storing the slot.
+// All atomics are seq_cst so no interleaving can release a ref twice or
+// free the task while a callback or the starter still touches it.
+struct UbrTimerTask {
+ butil::atomic<UbrTimerId>* slot;
+ butil::atomic<bthread_timer_t> id;
+ void* (*cb)(void*, uint64_t);
+ void* arg;
+ uint64_t gen; // opaque, passed back to cb
+ UbrTimerBackoffFn backoff;
+ uint64_t interval_us; // timer thread only
+ bool periodic;
+ butil::atomic<int> state; // kStarting/kScheduled/kDead
+ butil::atomic<bool> stopped;
+ butil::atomic<int> ref;
+ butil::atomic<bool> join_pending; // a DelAndWait is waiting
+ butil::atomic<bool> done; // refs hit zero, joiner frees
+};
+
+namespace {
+
+void ReleaseRef(UbrTimerTask* task) {
+ if (task->ref.fetch_sub(1) == 1) {
+ if (task->join_pending.load()) {
+ task->done.store(true); // joiner frees the task
+ } else {
+ delete task;
}
- return UBRING_ERR;
- }
- g_timer_module_initialized = 1;
- return UBRING_OK;
-}
-
-void *UnifiedCallback(void *args) {
- TimerFdCtx *ctx = (TimerFdCtx *)args;
- if (pthread_spin_lock(&ctx->spin_lock) != 0) {
- return nullptr;
- }
-
- if (ctx->status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&ctx->spin_lock);
- return nullptr;
- }
-
- void *(*cb)(void *) = ctx->cb;
- void *cb_args = ctx->args;
- uint32_t fd = ctx->fd;
- int is_periodical = ctx->periodical;
- ctx->status = TIMER_CONTEXT_CALLBACK_ONGOING;
-
- pthread_spin_unlock(&ctx->spin_lock);
-
- cb(cb_args);
-
- if (!is_periodical) {
- DeleteTimerInner(fd);
}
- return nullptr;
}
-void *TimerEpoll(void *args) {
- UNREFERENCE_PARAM(args);
-#if defined(OS_LINUX)
- struct epoll_event ready_events[MAX_TIMER];
-#elif defined(OS_MACOSX)
- struct kevent ready_events[MAX_TIMER];
-#endif
+void UbrTimerOnFire(void* p) {
+ UbrTimerTask* task = (UbrTimerTask*)p;
- while (1) {
- if (g_timer_module_initialized <= 0) {
- LOG(ERROR) << "The Timer module is not initialized.";
- break;
+ if (task->periodic) {
+ if (!task->stopped.load()) {
+ task->cb(task->arg, task->gen);
}
-
-#if defined(OS_LINUX)
- int32_t ready_num = epoll_wait(g_epoll_fd, ready_events, MAX_TIMER,
- TIMER_EPOLL_WAIT_TIMEOUT);
-#elif defined(OS_MACOSX)
- struct timespec timeout = {0, TIMER_EPOLL_WAIT_TIMEOUT * 1000000};
- int32_t ready_num = kevent(g_epoll_fd, nullptr, 0, ready_events,
MAX_TIMER, &timeout);
-#endif
-
- if (UNLIKELY(ready_num == -1)) {
- errno_t err = errno;
- if (err == EINTR) {
- LOG_EVERY_SECOND(WARNING) << "Epoll/Kqueue wait was
interrupted. errno=" << err;
- continue;
- } else if (err == EBADF) {
- LOG(WARNING) << "The Timer module is destroyed.";
- break;
+ // Claim the next schedule's ref before re-reading `stopped' so a
+ // racing delete can neither free the task nor orphan a re-arm.
+ task->ref.fetch_add(1);
+ if (task->stopped.load()) {
+ ReleaseRef(task);
+ } else {
+ uint64_t interval = task->interval_us;
+ if (task->backoff != nullptr) {
+ interval = task->backoff(task->arg, interval);
+ task->interval_us = interval;
}
- LOG(ERROR) << "Epoll/Kqueue wait internal error. errno=" << err;
- break;
- }
-
- for (int32_t i = 0; i < ready_num; i++) {
-#if defined(OS_LINUX)
- struct epoll_event *event = &ready_events[i];
- int32_t timer_fd = event->data.fd;
-#elif defined(OS_MACOSX)
- struct kevent *event = &ready_events[i];
- int32_t timer_fd = event->ident;
-#endif
-
- uint64_t exp = 0;
- if (read(timer_fd, &exp, sizeof(exp)) < 0) {
- if (errno != EBADF) {
- LOG(ERROR) << "Failed to read timerfd=" << timer_fd << "
errno=" << errno;
+ bthread_timer_t id = 0;
+ if (bthread_timer_add(
+ &id, butil::microseconds_from_now((int64_t)interval),
+ UbrTimerOnFire, task) == 0) {
+ task->id.store(id);
+ if (task->stopped.load() && bthread_timer_del(id) == 0) {
+ ReleaseRef(task);
}
- continue;
- }
- if (TimerFdCtxValidate((uint32_t)timer_fd) != UBRING_OK) {
- continue;
- }
-
- RETURN_CODE ret = ExecuteCallback(timer_fd);
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed execute callback ret=" << ret;
- DeleteTimerInner((uint32_t)timer_fd);
- continue;
+ } else {
+ LOG(ERROR) << "Fail to re-arm ubring timer";
+ ReleaseRef(task);
}
}
- }
- return nullptr;
-}
-
-void DeleteTimerSafe(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
+ ReleaseRef(task);
return;
}
- if (pthread_spin_lock(&g_timer_fd_ctx_map[fd].spin_lock) != 0) {
- return;
+ // One-shot: exit the handle slot first -- after this the wrapper never
+ // touches the storage again, so the callback may release the object
+ // that holds it. Whether the callback runs is decided solely by this
+ // slot competition: every UbrTimerDel that wants the callback
+ // suppressed has to win this exchange first, so owned==true guarantees
+ // no UbrTimerDel is pending. Do not consult `stopped' here: its store
+ // (del thread) and this load (timer thread) are separated by the slot
+ // RMW and seq_cst does not order the store-buffer case -- ownership of
+ // the slot is the single arbiter.
+ UbrTimerId expected = task;
+ const bool owned = task->slot->compare_exchange_strong(expected, nullptr);
+ if (owned) {
+ task->cb(task->arg, task->gen);
+ }
+ ReleaseRef(task); // schedule
+ if (owned) {
+ ReleaseRef(task); // owner
}
-
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
- return;
- }
-
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].periodical = 0;
- g_timer_fd_ctx_map[fd].fd = 0;
-
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
-
-#if defined(OS_LINUX)
- epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, (int)fd, nullptr);
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr);
-#endif
-
- uint64_t exp = 0;
- read((int)fd, &exp, sizeof(exp));
-
- close((int)fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
}
-void DeleteTimer(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
- LOG(WARNING) << "The timer is not initialized.";
- return;
- }
-
- g_timer_fd_ctx_map[fd].periodical = 0;
+UbrTimerTask* TakeOutTask(butil::atomic<UbrTimerId>* slot) {
+ return slot->exchange(nullptr);
}
-int32_t TimerStart(const itimerspec *time, void *(*cb)(void *), void *args) {
- if (g_epoll_fd == -1) {
- LOG(ERROR) << "Timer epoll/kqueue encountered internal error.";
- return -1;
- }
-
-#if defined(OS_LINUX)
- int timer_fd = timerfd_create(CLOCK_MONOTONIC, 0);
-#elif defined(OS_MACOSX)
- int timer_fd = timerfd_create_macosx(CLOCK_MONOTONIC, 0);
-#endif
-
- if (UNLIKELY(timer_fd >= (int)g_max_system_fd || timer_fd == -1)) {
- LOG(ERROR) << "Failed to create timerfd=" << timer_fd << " errno=" <<
errno;
- return -1;
+RETURN_CODE TimerStartInternal(butil::atomic<UbrTimerId>* slot, uint64_t
delay_us,
+ uint64_t interval_us, void* (*cb)(void*,
uint64_t),
+ void* arg, uint64_t gen,
+ UbrTimerBackoffFn backoff) {
+ if (BAIDU_UNLIKELY(slot == nullptr || cb == nullptr)) {
+ LOG(ERROR) << "Ubr timer start invalid argument, slot=" << slot;
+ return UBRING_ERR;
}
- g_timer_fd_ctx_map[timer_fd].status = TIMER_CONTEXT_EPOLL_WAITING;
- g_timer_fd_ctx_map[timer_fd].cb = cb;
- g_timer_fd_ctx_map[timer_fd].args = args;
- g_timer_fd_ctx_map[timer_fd].fd = (uint32_t)timer_fd;
-
- if (LIKELY(time->it_interval.tv_sec > 0 || time->it_interval.tv_nsec > 0))
{
- g_timer_fd_ctx_map[timer_fd].periodical = 1;
+ UbrTimerTask* task = new (std::nothrow) UbrTimerTask();
+ if (BAIDU_UNLIKELY(task == nullptr)) {
+ LOG(ERROR) << "Fail to malloc ubring timer task.";
+ return UBRING_ERR;
}
-
-#if defined(OS_LINUX)
- struct epoll_event event = {
- .events = EPOLLIN,
- .data = {.fd = timer_fd}
- };
-
- int32_t ret = epoll_ctl(g_epoll_fd, EPOLL_CTL_ADD, timer_fd, &event);
-#elif defined(OS_MACOSX)
- struct kevent event;
- uint64_t timeout_nsec = time->it_value.tv_sec * 1000000000ULL +
time->it_value.tv_nsec;
- uint64_t interval_nsec = time->it_interval.tv_sec * 1000000000ULL +
time->it_interval.tv_nsec;
- EV_SET(&event, timer_fd, EVFILT_TIMER, EV_ADD | EV_ENABLE, 0,
- timeout_nsec / 1000000, nullptr);
- int32_t ret = kevent(g_epoll_fd, &event, 1, nullptr, 0, nullptr);
-#endif
-
- if (UNLIKELY(ret != 0)) {
- CloseTimerFd(timer_fd);
- LOG(ERROR) << "Failed to add event to epoll/kqueue. errno=" << errno;
- return -1;
+ task->slot = slot;
+ task->id.store(0);
+ task->cb = cb;
+ task->arg = arg;
+ task->gen = gen;
+ task->backoff = backoff;
+ task->interval_us = interval_us;
+ task->periodic = (interval_us > 0);
+ task->state.store(kStarting);
+ task->stopped.store(false);
+ task->ref.store(3); // owner + schedule + starter
+ task->join_pending.store(false);
+ task->done.store(false);
+
+ // Publish the real task before scheduling so a delete or a DelAndWait
+ // racing the start always has an object to act on or wait for.
+ UbrTimerId expected = nullptr;
+ if (!slot->compare_exchange_strong(expected, task)) {
+ LOG(ERROR) << "Ubr timer start refused, slot already occupied";
+ delete task; // never published
+ return UBRING_ERR;
}
- std::atomic_fetch_add(&g_total_timer_num, 1U);
-
-#if defined(OS_LINUX)
- ret = timerfd_settime(timer_fd, 0, time, nullptr);
-#elif defined(OS_MACOSX)
- ret = timerfd_settime_macosx(timer_fd, 0, time, nullptr);
-#endif
-
- if (UNLIKELY(ret != 0)) {
-#if defined(OS_LINUX)
- if (epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, timer_fd, nullptr) != 0) {
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, timer_fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- if (kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr) != 0) {
-#endif
- LOG(ERROR) << "Failed to delete the timer fd=" << timer_fd << "
with errno=" << errno;
+ bthread_timer_t id = 0;
+ if (BAIDU_UNLIKELY(bthread_timer_add(
+ &id, butil::microseconds_from_now((int64_t)delay_us),
+ UbrTimerOnFire, task) != 0)) {
+ LOG(ERROR) << "Fail to add ubring timer";
+ task->state.store(kDead); // wake DelAndWait waiters
+ expected = task;
+ const bool owned = slot->compare_exchange_strong(expected, nullptr);
+ ReleaseRef(task); // schedule, never ran
+ if (owned) {
+ ReleaseRef(task); // owner
}
- CloseTimerFd(timer_fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
- LOG(ERROR) << "Failed to set timer";
- return -1;
+ ReleaseRef(task); // starter
+ return UBRING_ERR;
}
-
- return timer_fd;
+ // A zero-delay task may have fired and re-armed already; keep a newer
+ // id if so.
+ bthread_timer_t expected_id = 0;
+ task->id.compare_exchange_strong(expected_id, id);
+ task->state.store(kScheduled);
+ // No post-add stopped check here: a UbrTimerDel racing the start
+ // returns 1 without consuming the per-task resources, and the armed
+ // timer must fire so that OnFire settles the ownership protocol.
+ ReleaseRef(task); // starter
+ return UBRING_OK;
}
-uint32_t GetActiveTimerNum(void) {
- return std::atomic_load(&g_total_timer_num);
-}
+} // namespace
-void CloseTimerFd(int fd) {
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].fd = 0;
- g_timer_fd_ctx_map[fd].periodical = 0;
- if (close((int)fd) != 0) {
- LOG(ERROR) << "Failed to close timer fd=" << fd << " errno=" << errno;
- return;
- }
+RETURN_CODE UbrTimerStart(butil::atomic<UbrTimerId>* slot, uint64_t delay_us,
+ uint64_t interval_us, void* (*cb)(void*, uint64_t),
+ void* arg, uint64_t gen, UbrTimerBackoffFn backoff) {
+ return TimerStartInternal(slot, delay_us, interval_us, cb, arg, gen,
backoff);
}
-void TimerModuleDestroy(void) {
- uint32_t max_fd = g_max_system_fd;
- if (g_timer_fd_ctx_map) {
- for (uint32_t fd = 0; fd < max_fd; fd++) {
- if (g_timer_fd_ctx_map[fd].status != TIMER_CONTEXT_NOT_USING) {
- DeleteTimerSafe(fd);
- }
- }
- }
- close(g_epoll_fd);
- g_epoll_fd = -1;
- g_total_timer_num = 0;
- g_timer_module_initialized = 0;
- int32_t ret = pthread_join(g_epoll_execute_thread, nullptr);
- if (ret != EOK) {
- LOG(ERROR) << "Failed to join pthread, during destroying timer module.
ret=" << ret;
- return;
- }
+int UbrTimerDel(butil::atomic<UbrTimerId>* slot) {
+ if (slot == nullptr) {
+ return 1;
+ }
+ // Take the ownership of the slot first: after this exchange every
+ // dereference below is safe (the task cannot be freed while we hold
+ // the owner reference the slot used to anchor).
+ UbrTimerTask* task = TakeOutTask(slot);
+ if (task == nullptr) {
+ return 1; // fired and cleared its slot (callback side consumed)
+ // or another del won the exchange (it consumes)
+ }
+ task->stopped.store(true); // meaningful for periodic
only
+ // A start still in flight cannot be cancelled nor dispatched yet; wait
+ // for the starter to settle the fate (kScheduled/kDead). Bounded: the
+ // starter stores the state before taking any lock our caller holds.
+ while (task->state.load() == kStarting) {
+ bthread_usleep(1000);
+ }
+ if (task->state.load() == kDead) {
+ ReleaseRef(task); // owner; schedule/starter are
+ return 1; // settled by the kDead path
+ }
+ bthread_timer_t id = task->id.load();
+ if (id != 0 && bthread_timer_del(id) == 0) {
+ ReleaseRef(task); // schedule: cancelled before
dispatch
+ } // ==1: dispatched, OnFire
(owned==false)
+ // releases it
+ ReleaseRef(task); // owner
+ return 0; // This call won the slot competition. For a one-shot
timer,
+ // the callback will not run. For a periodic timer, future
+ // rearming is stopped, but an already dispatched or
running
+ // callback may still complete.
}
-RETURN_CODE TimerFdCtxValidate(uint32_t fd) {
- if (fd >= g_max_system_fd) {
- LOG(ERROR) << "TimerFd=" << fd << " is out of range=" <<
g_max_system_fd;
- return UBRING_ERR;
- }
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- LOG(ERROR) << "TimerFd=" << fd << " has wrong status=" <<
g_timer_fd_ctx_map[fd].status;
- return UBRING_ERR;
+void UbrTimerDelAndWait(butil::atomic<UbrTimerId>* slot) {
+ if (slot == nullptr) {
+ return;
}
- if (g_timer_fd_ctx_map[fd].cb == nullptr) {
- LOG(ERROR) << "The callback is not set.";
- return UBRING_ERR;
+ UbrTimerTask* task = TakeOutTask(slot);
+ if (task == nullptr) {
+ return;
}
Review Comment:
When another caller has already won the slot exchange (for example, a timer
callback calling `UbrTimerDel` on itself), this returns immediately even though
that task may still be executing. Teardown then proceeds to clear the `UbrTrx`
or `ShmList` while the callback still dereferences it, defeating the promised
wait semantics and allowing a use-after-free. The task being removed from the
slot needs a separate completion/join path that `UbrTimerDelAndWait` can
observe when it loses the slot race.
##########
src/brpc/ubshm/timer/timer_mgr.cpp:
##########
@@ -15,454 +15,258 @@
// specific language governing permissions and limitations
// under the License.
-#define _GNU_SOURCE
-#include <pthread.h>
-#include <sched.h>
-#include <errno.h>
-#include <stdio.h>
-#include <stdlib.h>
-#include <unistd.h>
-#include <atomic>
-#include <sys/resource.h>
+#include <new>
+#include "bthread/bthread.h" // bthread_usleep
+#include "bthread/unstable.h" // bthread_timer_add/del
+#include "butil/atomicops.h"
+#include "butil/time.h"
#include "brpc/ubshm/timer/timer_mgr.h"
namespace brpc {
namespace ubring {
-int32_t g_epoll_fd = -1;
-std::atomic<uint32_t> g_total_timer_num(0);
-TimerFdCtx *g_timer_fd_ctx_map = nullptr;
-uint32_t g_max_system_fd = 0;
-static pthread_t g_epoll_execute_thread = 0;
-static int32_t g_timer_module_initialized = 0;
-
-#if defined(OS_MACOSX)
-static int timerfd_create_macosx(int clockid, int flags);
-static int timerfd_settime_macosx(int fd, int flags,
- const itimerspec *new_value,
- itimerspec *old_value);
-#endif
-
-static RETURN_CODE DeleteTimerInner(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
- return UBRING_OK;
- }
-
- if (pthread_spin_lock(&g_timer_fd_ctx_map[fd].spin_lock) != 0) {
- return UBRING_ERR;
- }
-
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
- return UBRING_OK;
- }
-
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].periodical = 0;
- g_timer_fd_ctx_map[fd].fd = 0;
-
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
-
-#if defined(OS_LINUX)
- epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, (int)fd, nullptr);
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr);
-#endif
-
- uint64_t exp = 0;
- read((int)fd, &exp, sizeof(exp));
-
- close((int)fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
- return UBRING_OK;
-}
-
-static RETURN_CODE StartTimeEpoll(void) {
-#if defined(OS_LINUX)
- g_epoll_fd = epoll_create1(0);
-#elif defined(OS_MACOSX)
- g_epoll_fd = kqueue();
-#endif
- if (UNLIKELY(g_epoll_fd == -1)) {
- LOG(ERROR) << "Failed to create epoll/kqueue. errno=" << errno;
- return UBRING_ERR;
- }
-
- int ret = pthread_create(&g_epoll_execute_thread, nullptr, TimerEpoll,
nullptr);
- if (UNLIKELY(ret != 0)) {
- LOG(ERROR) << "Failed to create thread err=" << ret;
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerSpinLocksInit(void) {
- if (g_timer_fd_ctx_map == nullptr) {
- LOG(ERROR) << "Timer module is not fully initialized.";
- return UBRING_ERR;
- }
-
- for (uint32_t fd = 0; fd < g_max_system_fd; fd++) {
- int ret = pthread_spin_init(&g_timer_fd_ctx_map[fd].spin_lock,
- PTHREAD_PROCESS_PRIVATE);
- if (ret != EOK) {
- LOG(ERROR) << "Failed to initialize spin lock for fd=" << fd;
- for (uint32_t cleanup_fd = 0; cleanup_fd < fd; cleanup_fd++) {
-
pthread_spin_destroy(&g_timer_fd_ctx_map[cleanup_fd].spin_lock);
- }
- return UBRING_ERR;
- }
- }
- return UBRING_OK;
-}
-
-static RETURN_CODE ExecuteCallback(int32_t timer_fd) {
- UnifiedCallback((void *)(&g_timer_fd_ctx_map[timer_fd]));
- return UBRING_OK;
-}
-
-static RETURN_CODE TimerCtxMapCompletion(void) {
- memset(g_timer_fd_ctx_map, 0, sizeof(TimerFdCtx) * g_max_system_fd);
-
- RETURN_CODE ret = TimerSpinLocksInit();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init spin locks for timer module.";
- return UBRING_ERR;
- }
- return UBRING_OK;
-}
-
-RETURN_CODE TimerInit(void) {
- if (g_timer_module_initialized > 0) {
- return UBRING_OK;
- }
-
- g_total_timer_num.store(0);
-
- struct rlimit rlim;
- if (getrlimit(RLIMIT_NOFILE, &rlim) != UBRING_OK) {
- LOG(ERROR) << "Failed to get fd";
- return UBRING_ERR;
- }
- g_max_system_fd = (uint32_t)rlim.rlim_cur;
-
- if (g_timer_fd_ctx_map == nullptr) {
- g_timer_fd_ctx_map = (TimerFdCtx *)malloc(sizeof(TimerFdCtx) *
g_max_system_fd);
- if (UNLIKELY(!g_timer_fd_ctx_map)) {
- LOG(ERROR) << "Fail to malloc space for timer modules. errno=%d",
errno;
- return UBRING_ERR;
- }
-
- RETURN_CODE ret = TimerCtxMapCompletion();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to init main data structure of Time Module.
ret=" << ret;
- free(g_timer_fd_ctx_map);
- g_timer_fd_ctx_map = nullptr;
- return UBRING_ERR;
- }
- }
-
- RETURN_CODE ret = StartTimeEpoll();
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed to start Timer Epoll. ret=" << ret;
- if (LIKELY(g_timer_fd_ctx_map != nullptr)) {
- FREE_PTR(g_timer_fd_ctx_map);
+namespace {
+
+enum UbrTimerState {
+ kStarting = 0, // published, not scheduled
yet
+ kScheduled = 1,
+ kDead = 2 // scheduling failed
+};
+
+} // namespace
+
+// Reference rules: one "owner" ref for the handle slot, one "schedule" ref
+// per pending/running bthread schedule, plus one ref held by the starter
+// until its post-schedule bookkeeping is done. The schedule ref is
+// consumed by the firing callback or by the deleter whose
+// bthread_timer_del returned 0 (cancelled before run); the owner ref is
+// consumed by whoever takes the task out of *slot -- a deleter, or the
+// one-shot firing callback itself, which exits the slot BEFORE running
+// the callback so that the callback may free the object storing the slot.
+// All atomics are seq_cst so no interleaving can release a ref twice or
+// free the task while a callback or the starter still touches it.
+struct UbrTimerTask {
+ butil::atomic<UbrTimerId>* slot;
+ butil::atomic<bthread_timer_t> id;
+ void* (*cb)(void*, uint64_t);
+ void* arg;
+ uint64_t gen; // opaque, passed back to cb
+ UbrTimerBackoffFn backoff;
+ uint64_t interval_us; // timer thread only
+ bool periodic;
+ butil::atomic<int> state; // kStarting/kScheduled/kDead
+ butil::atomic<bool> stopped;
+ butil::atomic<int> ref;
+ butil::atomic<bool> join_pending; // a DelAndWait is waiting
+ butil::atomic<bool> done; // refs hit zero, joiner frees
+};
+
+namespace {
+
+void ReleaseRef(UbrTimerTask* task) {
+ if (task->ref.fetch_sub(1) == 1) {
+ if (task->join_pending.load()) {
+ task->done.store(true); // joiner frees the task
+ } else {
+ delete task;
}
- return UBRING_ERR;
- }
- g_timer_module_initialized = 1;
- return UBRING_OK;
-}
-
-void *UnifiedCallback(void *args) {
- TimerFdCtx *ctx = (TimerFdCtx *)args;
- if (pthread_spin_lock(&ctx->spin_lock) != 0) {
- return nullptr;
- }
-
- if (ctx->status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&ctx->spin_lock);
- return nullptr;
- }
-
- void *(*cb)(void *) = ctx->cb;
- void *cb_args = ctx->args;
- uint32_t fd = ctx->fd;
- int is_periodical = ctx->periodical;
- ctx->status = TIMER_CONTEXT_CALLBACK_ONGOING;
-
- pthread_spin_unlock(&ctx->spin_lock);
-
- cb(cb_args);
-
- if (!is_periodical) {
- DeleteTimerInner(fd);
}
- return nullptr;
}
-void *TimerEpoll(void *args) {
- UNREFERENCE_PARAM(args);
-#if defined(OS_LINUX)
- struct epoll_event ready_events[MAX_TIMER];
-#elif defined(OS_MACOSX)
- struct kevent ready_events[MAX_TIMER];
-#endif
+void UbrTimerOnFire(void* p) {
+ UbrTimerTask* task = (UbrTimerTask*)p;
- while (1) {
- if (g_timer_module_initialized <= 0) {
- LOG(ERROR) << "The Timer module is not initialized.";
- break;
+ if (task->periodic) {
+ if (!task->stopped.load()) {
+ task->cb(task->arg, task->gen);
}
-
-#if defined(OS_LINUX)
- int32_t ready_num = epoll_wait(g_epoll_fd, ready_events, MAX_TIMER,
- TIMER_EPOLL_WAIT_TIMEOUT);
-#elif defined(OS_MACOSX)
- struct timespec timeout = {0, TIMER_EPOLL_WAIT_TIMEOUT * 1000000};
- int32_t ready_num = kevent(g_epoll_fd, nullptr, 0, ready_events,
MAX_TIMER, &timeout);
-#endif
-
- if (UNLIKELY(ready_num == -1)) {
- errno_t err = errno;
- if (err == EINTR) {
- LOG_EVERY_SECOND(WARNING) << "Epoll/Kqueue wait was
interrupted. errno=" << err;
- continue;
- } else if (err == EBADF) {
- LOG(WARNING) << "The Timer module is destroyed.";
- break;
+ // Claim the next schedule's ref before re-reading `stopped' so a
+ // racing delete can neither free the task nor orphan a re-arm.
+ task->ref.fetch_add(1);
+ if (task->stopped.load()) {
+ ReleaseRef(task);
+ } else {
+ uint64_t interval = task->interval_us;
+ if (task->backoff != nullptr) {
+ interval = task->backoff(task->arg, interval);
+ task->interval_us = interval;
}
- LOG(ERROR) << "Epoll/Kqueue wait internal error. errno=" << err;
- break;
- }
-
- for (int32_t i = 0; i < ready_num; i++) {
-#if defined(OS_LINUX)
- struct epoll_event *event = &ready_events[i];
- int32_t timer_fd = event->data.fd;
-#elif defined(OS_MACOSX)
- struct kevent *event = &ready_events[i];
- int32_t timer_fd = event->ident;
-#endif
-
- uint64_t exp = 0;
- if (read(timer_fd, &exp, sizeof(exp)) < 0) {
- if (errno != EBADF) {
- LOG(ERROR) << "Failed to read timerfd=" << timer_fd << "
errno=" << errno;
+ bthread_timer_t id = 0;
+ if (bthread_timer_add(
+ &id, butil::microseconds_from_now((int64_t)interval),
+ UbrTimerOnFire, task) == 0) {
+ task->id.store(id);
+ if (task->stopped.load() && bthread_timer_del(id) == 0) {
+ ReleaseRef(task);
}
- continue;
- }
- if (TimerFdCtxValidate((uint32_t)timer_fd) != UBRING_OK) {
- continue;
- }
-
- RETURN_CODE ret = ExecuteCallback(timer_fd);
- if (ret != UBRING_OK) {
- LOG(ERROR) << "Failed execute callback ret=" << ret;
- DeleteTimerInner((uint32_t)timer_fd);
- continue;
+ } else {
+ LOG(ERROR) << "Fail to re-arm ubring timer";
+ ReleaseRef(task);
}
}
- }
- return nullptr;
-}
-
-void DeleteTimerSafe(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
+ ReleaseRef(task);
return;
}
- if (pthread_spin_lock(&g_timer_fd_ctx_map[fd].spin_lock) != 0) {
- return;
+ // One-shot: exit the handle slot first -- after this the wrapper never
+ // touches the storage again, so the callback may release the object
+ // that holds it. Whether the callback runs is decided solely by this
+ // slot competition: every UbrTimerDel that wants the callback
+ // suppressed has to win this exchange first, so owned==true guarantees
+ // no UbrTimerDel is pending. Do not consult `stopped' here: its store
+ // (del thread) and this load (timer thread) are separated by the slot
+ // RMW and seq_cst does not order the store-buffer case -- ownership of
+ // the slot is the single arbiter.
+ UbrTimerId expected = task;
+ const bool owned = task->slot->compare_exchange_strong(expected, nullptr);
+ if (owned) {
+ task->cb(task->arg, task->gen);
+ }
+ ReleaseRef(task); // schedule
+ if (owned) {
+ ReleaseRef(task); // owner
}
-
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
- return;
- }
-
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].periodical = 0;
- g_timer_fd_ctx_map[fd].fd = 0;
-
- pthread_spin_unlock(&g_timer_fd_ctx_map[fd].spin_lock);
-
-#if defined(OS_LINUX)
- epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, (int)fd, nullptr);
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr);
-#endif
-
- uint64_t exp = 0;
- read((int)fd, &exp, sizeof(exp));
-
- close((int)fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
}
-void DeleteTimer(uint32_t fd) {
- if (g_timer_fd_ctx_map == nullptr) {
- LOG(WARNING) << "The timer is not initialized.";
- return;
- }
-
- g_timer_fd_ctx_map[fd].periodical = 0;
+UbrTimerTask* TakeOutTask(butil::atomic<UbrTimerId>* slot) {
+ return slot->exchange(nullptr);
}
-int32_t TimerStart(const itimerspec *time, void *(*cb)(void *), void *args) {
- if (g_epoll_fd == -1) {
- LOG(ERROR) << "Timer epoll/kqueue encountered internal error.";
- return -1;
- }
-
-#if defined(OS_LINUX)
- int timer_fd = timerfd_create(CLOCK_MONOTONIC, 0);
-#elif defined(OS_MACOSX)
- int timer_fd = timerfd_create_macosx(CLOCK_MONOTONIC, 0);
-#endif
-
- if (UNLIKELY(timer_fd >= (int)g_max_system_fd || timer_fd == -1)) {
- LOG(ERROR) << "Failed to create timerfd=" << timer_fd << " errno=" <<
errno;
- return -1;
+RETURN_CODE TimerStartInternal(butil::atomic<UbrTimerId>* slot, uint64_t
delay_us,
+ uint64_t interval_us, void* (*cb)(void*,
uint64_t),
+ void* arg, uint64_t gen,
+ UbrTimerBackoffFn backoff) {
+ if (BAIDU_UNLIKELY(slot == nullptr || cb == nullptr)) {
+ LOG(ERROR) << "Ubr timer start invalid argument, slot=" << slot;
+ return UBRING_ERR;
}
- g_timer_fd_ctx_map[timer_fd].status = TIMER_CONTEXT_EPOLL_WAITING;
- g_timer_fd_ctx_map[timer_fd].cb = cb;
- g_timer_fd_ctx_map[timer_fd].args = args;
- g_timer_fd_ctx_map[timer_fd].fd = (uint32_t)timer_fd;
-
- if (LIKELY(time->it_interval.tv_sec > 0 || time->it_interval.tv_nsec > 0))
{
- g_timer_fd_ctx_map[timer_fd].periodical = 1;
+ UbrTimerTask* task = new (std::nothrow) UbrTimerTask();
+ if (BAIDU_UNLIKELY(task == nullptr)) {
+ LOG(ERROR) << "Fail to malloc ubring timer task.";
+ return UBRING_ERR;
}
-
-#if defined(OS_LINUX)
- struct epoll_event event = {
- .events = EPOLLIN,
- .data = {.fd = timer_fd}
- };
-
- int32_t ret = epoll_ctl(g_epoll_fd, EPOLL_CTL_ADD, timer_fd, &event);
-#elif defined(OS_MACOSX)
- struct kevent event;
- uint64_t timeout_nsec = time->it_value.tv_sec * 1000000000ULL +
time->it_value.tv_nsec;
- uint64_t interval_nsec = time->it_interval.tv_sec * 1000000000ULL +
time->it_interval.tv_nsec;
- EV_SET(&event, timer_fd, EVFILT_TIMER, EV_ADD | EV_ENABLE, 0,
- timeout_nsec / 1000000, nullptr);
- int32_t ret = kevent(g_epoll_fd, &event, 1, nullptr, 0, nullptr);
-#endif
-
- if (UNLIKELY(ret != 0)) {
- CloseTimerFd(timer_fd);
- LOG(ERROR) << "Failed to add event to epoll/kqueue. errno=" << errno;
- return -1;
+ task->slot = slot;
+ task->id.store(0);
+ task->cb = cb;
+ task->arg = arg;
+ task->gen = gen;
+ task->backoff = backoff;
+ task->interval_us = interval_us;
+ task->periodic = (interval_us > 0);
+ task->state.store(kStarting);
+ task->stopped.store(false);
+ task->ref.store(3); // owner + schedule + starter
+ task->join_pending.store(false);
+ task->done.store(false);
+
+ // Publish the real task before scheduling so a delete or a DelAndWait
+ // racing the start always has an object to act on or wait for.
+ UbrTimerId expected = nullptr;
+ if (!slot->compare_exchange_strong(expected, task)) {
+ LOG(ERROR) << "Ubr timer start refused, slot already occupied";
+ delete task; // never published
+ return UBRING_ERR;
}
- std::atomic_fetch_add(&g_total_timer_num, 1U);
-
-#if defined(OS_LINUX)
- ret = timerfd_settime(timer_fd, 0, time, nullptr);
-#elif defined(OS_MACOSX)
- ret = timerfd_settime_macosx(timer_fd, 0, time, nullptr);
-#endif
-
- if (UNLIKELY(ret != 0)) {
-#if defined(OS_LINUX)
- if (epoll_ctl(g_epoll_fd, EPOLL_CTL_DEL, timer_fd, nullptr) != 0) {
-#elif defined(OS_MACOSX)
- struct kevent evt;
- EV_SET(&evt, timer_fd, EVFILT_TIMER, EV_DELETE, 0, 0, nullptr);
- if (kevent(g_epoll_fd, &evt, 1, nullptr, 0, nullptr) != 0) {
-#endif
- LOG(ERROR) << "Failed to delete the timer fd=" << timer_fd << "
with errno=" << errno;
+ bthread_timer_t id = 0;
+ if (BAIDU_UNLIKELY(bthread_timer_add(
+ &id, butil::microseconds_from_now((int64_t)delay_us),
+ UbrTimerOnFire, task) != 0)) {
+ LOG(ERROR) << "Fail to add ubring timer";
+ task->state.store(kDead); // wake DelAndWait waiters
+ expected = task;
+ const bool owned = slot->compare_exchange_strong(expected, nullptr);
+ ReleaseRef(task); // schedule, never ran
+ if (owned) {
+ ReleaseRef(task); // owner
}
- CloseTimerFd(timer_fd);
- std::atomic_fetch_sub(&g_total_timer_num, 1U);
- LOG(ERROR) << "Failed to set timer";
- return -1;
+ ReleaseRef(task); // starter
+ return UBRING_ERR;
}
-
- return timer_fd;
+ // A zero-delay task may have fired and re-armed already; keep a newer
+ // id if so.
+ bthread_timer_t expected_id = 0;
+ task->id.compare_exchange_strong(expected_id, id);
+ task->state.store(kScheduled);
+ // No post-add stopped check here: a UbrTimerDel racing the start
+ // returns 1 without consuming the per-task resources, and the armed
+ // timer must fire so that OnFire settles the ownership protocol.
+ ReleaseRef(task); // starter
+ return UBRING_OK;
}
-uint32_t GetActiveTimerNum(void) {
- return std::atomic_load(&g_total_timer_num);
-}
+} // namespace
-void CloseTimerFd(int fd) {
- g_timer_fd_ctx_map[fd].cb = nullptr;
- g_timer_fd_ctx_map[fd].args = nullptr;
- g_timer_fd_ctx_map[fd].status = TIMER_CONTEXT_NOT_USING;
- g_timer_fd_ctx_map[fd].fd = 0;
- g_timer_fd_ctx_map[fd].periodical = 0;
- if (close((int)fd) != 0) {
- LOG(ERROR) << "Failed to close timer fd=" << fd << " errno=" << errno;
- return;
- }
+RETURN_CODE UbrTimerStart(butil::atomic<UbrTimerId>* slot, uint64_t delay_us,
+ uint64_t interval_us, void* (*cb)(void*, uint64_t),
+ void* arg, uint64_t gen, UbrTimerBackoffFn backoff) {
+ return TimerStartInternal(slot, delay_us, interval_us, cb, arg, gen,
backoff);
}
-void TimerModuleDestroy(void) {
- uint32_t max_fd = g_max_system_fd;
- if (g_timer_fd_ctx_map) {
- for (uint32_t fd = 0; fd < max_fd; fd++) {
- if (g_timer_fd_ctx_map[fd].status != TIMER_CONTEXT_NOT_USING) {
- DeleteTimerSafe(fd);
- }
- }
- }
- close(g_epoll_fd);
- g_epoll_fd = -1;
- g_total_timer_num = 0;
- g_timer_module_initialized = 0;
- int32_t ret = pthread_join(g_epoll_execute_thread, nullptr);
- if (ret != EOK) {
- LOG(ERROR) << "Failed to join pthread, during destroying timer module.
ret=" << ret;
- return;
- }
+int UbrTimerDel(butil::atomic<UbrTimerId>* slot) {
+ if (slot == nullptr) {
+ return 1;
+ }
+ // Take the ownership of the slot first: after this exchange every
+ // dereference below is safe (the task cannot be freed while we hold
+ // the owner reference the slot used to anchor).
+ UbrTimerTask* task = TakeOutTask(slot);
+ if (task == nullptr) {
+ return 1; // fired and cleared its slot (callback side consumed)
+ // or another del won the exchange (it consumes)
+ }
+ task->stopped.store(true); // meaningful for periodic
only
+ // A start still in flight cannot be cancelled nor dispatched yet; wait
+ // for the starter to settle the fate (kScheduled/kDead). Bounded: the
+ // starter stores the state before taking any lock our caller holds.
+ while (task->state.load() == kStarting) {
+ bthread_usleep(1000);
+ }
+ if (task->state.load() == kDead) {
+ ReleaseRef(task); // owner; schedule/starter are
+ return 1; // settled by the kDead path
+ }
+ bthread_timer_t id = task->id.load();
+ if (id != 0 && bthread_timer_del(id) == 0) {
+ ReleaseRef(task); // schedule: cancelled before
dispatch
+ } // ==1: dispatched, OnFire
(owned==false)
+ // releases it
+ ReleaseRef(task); // owner
+ return 0; // This call won the slot competition. For a one-shot
timer,
+ // the callback will not run. For a periodic timer, future
+ // rearming is stopped, but an already dispatched or
running
+ // callback may still complete.
}
-RETURN_CODE TimerFdCtxValidate(uint32_t fd) {
- if (fd >= g_max_system_fd) {
- LOG(ERROR) << "TimerFd=" << fd << " is out of range=" <<
g_max_system_fd;
- return UBRING_ERR;
- }
- if (g_timer_fd_ctx_map[fd].status == TIMER_CONTEXT_NOT_USING) {
- LOG(ERROR) << "TimerFd=" << fd << " has wrong status=" <<
g_timer_fd_ctx_map[fd].status;
- return UBRING_ERR;
+void UbrTimerDelAndWait(butil::atomic<UbrTimerId>* slot) {
Review Comment:
The new timer facade is a custom ref-counted state machine, but the UBRing
tests added here only verify flag metadata; no test exercises one-shot
cancellation, periodic rearming/backoff, concurrent delete versus callback, or
`UbrTimerDelAndWait`. These lifecycle paths are the purpose of this change and
regressions can corrupt pooled trx lifetime, so add focused unit tests for the
timer wrapper (including callback self-delete and teardown races).
--
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]