169 using Clock = std::chrono::steady_clock;
173 static constexpr bool TracksTime = Config.maxIdleTimeMs > 0 || Config.maxLifetimeMs > 0;
184 std::unique_ptr<DataMapper> mapper;
185 Clock::time_point createdAt {};
186 Clock::time_point idleSince {};
200 _entry { std::move(entry) },
212 _entry { std::move(other._entry) },
213 _pool { other._pool }
216 PooledDataMapper& operator=(PooledDataMapper
const&) =
delete;
217 PooledDataMapper& operator=(PooledDataMapper&&) =
delete;
218 ~PooledDataMapper() noexcept
227 return _entry.mapper.get();
235 return *_entry.mapper;
239 void ReturnToPool() noexcept
241 _pool.Return(std::move(_entry));
242 _entry.mapper =
nullptr;
260 static void DropAsyncBackend(DataMapper& dm)
noexcept
262 dm.Connection().DisableAsync();
267 [[nodiscard]] Clock::time_point NowIfTracking() const noexcept
269 if constexpr (!TracksTime)
281 [[nodiscard]] Entry MakeEntry()
const
283 auto const now = NowIfTracking();
284 auto mapper = std::make_unique<DataMapper>();
288 if constexpr (Config.preparedStatementCacheCapacity != 0)
289 mapper->Connection().SetPreparedStatementCacheCapacity(Config.preparedStatementCacheCapacity);
290 return Entry { std::move(mapper), now, now };
302 [[nodiscard]]
bool IsUsable(Entry
const& entry, Clock::time_point now)
const noexcept
304 if constexpr (Config.maxLifetimeMs > 0)
306 if (now - entry.createdAt >= Config.MaxLifetime())
309 if constexpr (Config.maxIdleTimeMs > 0)
311 if (now - entry.idleSince >= Config.MaxIdleTime())
316 if (!entry.mapper->Connection().IsAlive())
331 [[nodiscard]] Entry TakeIdleLocked(std::vector<Entry>& retired)
333 auto const now = NowIfTracking();
334 while (!_idleDataMappers.empty())
336 auto entry = std::move(_idleDataMappers.back());
337 _idleDataMappers.pop_back();
338 if (IsUsable(entry, now))
340 retired.push_back(std::move(entry));
351 [[nodiscard]]
static bool IsPastLifetime([[maybe_unused]] Entry
const& entry,
352 [[maybe_unused]] Clock::time_point now)
noexcept
354 if constexpr (Config.maxLifetimeMs > 0)
355 return now - entry.createdAt >= Config.MaxLifetime();
361 void Return(Entry entry)
noexcept
364 DropAsyncBackend(*entry.mapper);
365 auto const now = NowIfTracking();
366 if (IsPastLifetime(entry, now))
370 LIGHTWEIGHT_STATS_POOL_RELEASE(
true);
373 entry.idleSince = now;
375 std::scoped_lock lock(_mutex);
376 _idleDataMappers.push_back(std::move(entry));
377 LIGHTWEIGHT_STATS_POOL_RELEASE(
false);
378 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
383 void Return(Entry entry)
noexcept
386 DropAsyncBackend(*entry.mapper);
388 std::shared_ptr<WaiterNode> toResume;
390 std::scoped_lock
const lock(_mutex);
391 toResume = ReturnLocked(std::move(entry), retired);
395 toResume->resume->Resume(toResume->handle);
405 [[nodiscard]] Entry AcquireReadyLocked(std::vector<Entry>& retired)
408 if (
auto entry = TakeIdleLocked(retired); entry.mapper)
414 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
true,
false);
415 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
418 if (_checkedOut < Config.maxSize)
422 auto fresh = MakeEntry();
424 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
false,
false);
425 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
442 std::shared_ptr<WaiterNode> ReturnLocked(Entry entry, Entry& retired)
noexcept
445 while (!_waiters.empty())
447 auto node = _waiters.front();
448 _waiters.pop_front();
451 if (node->state != WaiterNode::State::Parked)
462 LIGHTWEIGHT_STATS_POOL_RELEASE(
false);
463 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
464 node->state = WaiterNode::State::Fulfilled;
465 node->entry = std::move(entry);
466 if (node->kind == WaiterNode::Kind::Async)
468 node->cv.notify_one();
474 auto const now = NowIfTracking();
475 if (IsPastLifetime(entry, now))
477 retired = std::move(entry);
480 LIGHTWEIGHT_STATS_POOL_RELEASE(
true);
481 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
484 entry.idleSince = now;
486 _idleDataMappers.push_back(std::move(entry));
487 LIGHTWEIGHT_STATS_POOL_RELEASE(
false);
488 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
493 void Return(Entry entry)
noexcept
496 DropAsyncBackend(*entry.mapper);
497 auto const now = NowIfTracking();
498 if (IsPastLifetime(entry, now))
502 LIGHTWEIGHT_STATS_POOL_RELEASE(
true);
505 entry.idleSince = now;
506 std::scoped_lock lock(_mutex);
507 if (_idleDataMappers.size() < Config.maxSize)
510 _idleDataMappers.push_back(std::move(entry));
511 LIGHTWEIGHT_STATS_POOL_RELEASE(
false);
517 LIGHTWEIGHT_STATS_POOL_RELEASE(
true);
519 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
527 _idleDataMappers.reserve(Config.initialSize);
528 for ([[maybe_unused]]
auto const _: std::views::iota(0U, Config.initialSize))
529 _idleDataMappers.push_back(MakeEntry());
541 if (!_waiters.empty())
543 "Pool destroyed while acquirers are still waiting on it (coroutines parked in AcquireAsync "
544 "and/or threads blocked in Acquire); the pool must outlive every acquirer (drive each "
545 "AcquireAsync task to completion or destroy it first, and never destroy the pool while a "
546 "thread is blocked in Acquire). This is undefined behavior.");
547 assert(_waiters.empty() &&
"Pool destroyed while acquirers are still waiting on it");
551 Pool& operator=(
Pool const&) =
delete;
564 std::vector<Entry> retired;
565 [[maybe_unused]]
auto const acquireStartedAt = std::chrono::steady_clock::now();
566 std::unique_lock lock(_mutex);
567 if (
auto entry = AcquireReadyLocked(retired); entry.mapper)
572 auto node = std::make_shared<WaiterNode>(WaiterNode::Kind::Sync);
573 _waiters.push_back(node);
574 node->cv.wait(lock, [&node] {
return node->state == WaiterNode::State::Fulfilled; });
576 LIGHTWEIGHT_STATS_POOL_ACQUIRE(
577 std::chrono::duration_cast<std::chrono::microseconds>(std::chrono::steady_clock::now() - acquireStartedAt),
580 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
592 [[nodiscard]] std::expected<PooledDataMapper, PoolError>
Acquire(std::chrono::milliseconds timeout)
595 std::vector<Entry> retired;
596 [[maybe_unused]]
auto const acquireStartedAt = std::chrono::steady_clock::now();
597 std::unique_lock lock(_mutex);
598 if (
auto entry = AcquireReadyLocked(retired); entry.mapper)
601 auto node = std::make_shared<WaiterNode>(WaiterNode::Kind::Sync);
602 _waiters.push_back(node);
603 if (!node->cv.wait_for(lock, timeout, [&node] { return node->state == WaiterNode::State::Fulfilled; }))
608 node->state = WaiterNode::State::Abandoned;
609 std::erase(_waiters, node);
614 LIGHTWEIGHT_STATS_POOL_ACQUIRE(
615 std::chrono::duration_cast<std::chrono::microseconds>(std::chrono::steady_clock::now() - acquireStartedAt),
618 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
628 std::vector<Entry> retired;
629 std::scoped_lock lock(_mutex);
630 auto entry = TakeIdleLocked(retired);
634 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
false,
false);
635 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
639 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
true,
false);
640 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(_idleDataMappers.size(), _checkedOut);
652 [[nodiscard]] std::expected<PooledDataMapper, PoolError>
Acquire([[maybe_unused]] std::chrono::milliseconds timeout)
671 return AcquireAsyncImpl(&dbWorkers, &resume);
690 _asyncDbWorkers = &dbWorkers;
691 _asyncResume = &resume;
704 if (!_asyncDbWorkers || !_asyncResume)
705 throw std::logic_error {
706 "Pool::AcquireAsync(): no async executors configured; call Pool::SetAsyncExecutors(...) first "
707 "or use the explicit AcquireAsync(dbWorkers, resume) overload."
709 return AcquireAsyncImpl(_asyncDbWorkers, _asyncResume);
719 void SetClock(std::function<Clock::time_point()> clock)
noexcept
721 _clock = std::move(clock);
724#if defined(BUILD_TESTS)
725 [[nodiscard]]
size_t IdleCount() noexcept
727 std::scoped_lock lock(_mutex);
728 return _idleDataMappers.size();
733 [[nodiscard]]
size_t WaiterCount() noexcept
735 std::scoped_lock lock(_mutex);
736 return _waiters.size();
751 enum class Kind : std::uint8_t
758 enum class State : std::uint8_t
766 State state = State::Parked;
770 std::coroutine_handle<> handle {};
771 Async::IResumeScheduler* resume =
nullptr;
775 std::condition_variable cv {};
777 explicit WaiterNode(Kind nodeKind)
noexcept:
788 struct AsyncAcquireAwaitable
791 Async::IResumeScheduler& resume;
793 std::shared_ptr<WaiterNode> node {};
794 std::vector<Entry> retired {};
797 std::chrono::steady_clock::time_point parkedAt {};
799 AsyncAcquireAwaitable(
Pool& poolRef, Async::IResumeScheduler& resumeRef)
noexcept:
805 AsyncAcquireAwaitable(AsyncAcquireAwaitable
const&) =
delete;
806 AsyncAcquireAwaitable& operator=(AsyncAcquireAwaitable
const&) =
delete;
807 AsyncAcquireAwaitable(AsyncAcquireAwaitable&&) =
delete;
808 AsyncAcquireAwaitable& operator=(AsyncAcquireAwaitable&&) =
delete;
821 ~AsyncAcquireAwaitable()
826 std::shared_ptr<WaiterNode> toResume;
828 std::scoped_lock
const lock(pool._mutex);
831 case WaiterNode::State::Parked:
834 node->state = WaiterNode::State::Abandoned;
835 std::erase(pool._waiters, node);
837 case WaiterNode::State::Fulfilled:
840 node->state = WaiterNode::State::Abandoned;
843 if (node->entry.mapper)
844 toResume = pool.ReturnLocked(std::move(node->entry), reclaimed);
847 case WaiterNode::State::Abandoned:
852 toResume->resume->Resume(toResume->handle);
855 [[nodiscard]]
bool await_ready() const noexcept
860 bool await_suspend(std::coroutine_handle<> handle)
862 std::scoped_lock
const lock(pool._mutex);
865 acquired = pool.TakeIdleLocked(retired);
871 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
true,
false);
872 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(pool._idleDataMappers.size(), pool._checkedOut);
880 if (pool._checkedOut >= Config.maxSize)
882 node = std::make_shared<WaiterNode>(WaiterNode::Kind::Async);
883 node->handle = handle;
884 node->resume = &resume;
885 pool._waiters.push_back(node);
886 parkedAt = std::chrono::steady_clock::now();
894 auto fresh = pool.MakeEntry();
897 acquired = std::move(fresh);
898 LIGHTWEIGHT_STATS_POOL_ACQUIRE(std::chrono::microseconds { 0 },
false,
false);
899 LIGHTWEIGHT_STATS_POOL_OCCUPANCY(pool._idleDataMappers.size(), pool._checkedOut);
903 Entry await_resume() noexcept
911 LIGHTWEIGHT_STATS_POOL_ACQUIRE(
912 std::chrono::duration_cast<std::chrono::microseconds>(std::chrono::steady_clock::now() - parkedAt),
915 return std::move(node->entry);
917 return std::move(acquired);
921 Async::Task<PooledDataMapper> AcquireAsyncImpl(Async::IExecutor* dbWorkers, Async::IResumeScheduler* resume)
923 auto entry =
co_await AsyncAcquireAwaitable { *
this, *resume };
927 auto pooled = PooledDataMapper(*
this, std::move(entry));
928 pooled->Connection().EnableAsync(*dbWorkers, *resume);
929 co_return std::move(pooled);
933 std::vector<Entry> _idleDataMappers;
934 size_t _checkedOut {};
941 std::function<Clock::time_point()> _clock {};
944 Async::IExecutor* _asyncDbWorkers =
nullptr;
945 Async::IResumeScheduler* _asyncResume =
nullptr;
948 std::deque<std::shared_ptr<WaiterNode>> _waiters;