/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
test/common/xenium/boundedqueueadaptor.hpp
148 строк
5 KB
Ant010ff
Demo code update, 2026.
08 фев 2026, 14:26
08 фев 2026, 14:26
1a5a363
Код
Авторство
О чём код?
/* xenium bounded queues (https://github.com/mpoeter/xenium) adaptor for FFPP library. Official repository: https://gitlab.com/ant010ff/ffpp Licensed under the MIT License <http://opensource.org/licenses/MIT>. SPDX-License-Identifier: MIT Copyright (c) 2021 - 2026 Anton Nasonov <ant010fff @ gmail . com>. */ #ifndef XENIUM_BOUNDED_QUEUE_ADAPTOR_HPP #define XENIUM_BOUNDED_QUEUE_ADAPTOR_HPP ////XeniumBoundedQueuePolicy////////////////////////////////////////////////////////////////////////////////////////////////////// struct XeniumBoundedQueuePolicy { struct MutexPolicy : ffpp::StripedTicketMutexPolicy { static size_t SharedSlots() { auto const nThreads = std::thread::hardware_concurrency(); if(nThreads < 2) return 1; //8 shared slots is a reasonable limit considering scalability limit of a Vyukov Bounded MPMC Queue algorithm. else return std::clamp(size_t(nThreads / 4), size_t(2), size_t(8)); } }; using MutexT = ffpp::StripedTicketMutex<ffpp::MutexFlags::Default, MutexPolicy>; using ResourcePolicyT = ffpp::DefaultResourcePolicy; using LockableT = ffpp::Lockable<MutexT>; static consteval auto Backoff() { return ffpp::JitteredBackoff<true, 6, 0, 0, 1, 5>(); } static constexpr auto QueueMode() { return ffpp::QueueFlags::Grow; } static size_t QueueCapacity() { return 2 * std::thread::hardware_concurrency(); } }; ////XeniumBoundedQueueAdaptor///////////////////////////////////////////////////////////////////////////////////////////////////// template< template <typename> typename TBoundedQueue, std::movable TValue, ffpp::Concepts::QueuePolicy TPolicy = XeniumBoundedQueuePolicy > class XeniumBoundedQueueAdaptor : public TPolicy::LockableT { public: using ValueT = TValue; using PolicyT = TPolicy; using ResourcePolicyT = PolicyT::ResourcePolicyT; using LockableT = PolicyT::LockableT; using UniqueLockT = LockableT::UniqueLockT; using SharedLockT = LockableT::SharedLockT; template<typename T> using AllocatorT = typename ResourcePolicyT::template AllocatorT<T>; static size_t constexpr c_nInvalidLimit = ResourcePolicyT::template InvalidValue<size_t>(); static consteval bool IsMode(ffpp::QueueFlags flg, bool bAll = false) { if(bAll) return (flg & PolicyT::QueueMode()) == flg; else return (flg & PolicyT::QueueMode()) != ffpp::QueueFlags::Undefined; } XeniumBoundedQueueAdaptor(size_t nLimit = c_nInvalidLimit) : nLimit_(nLimit) { } XeniumBoundedQueueAdaptor(XeniumBoundedQueueAdaptor const&) = delete; XeniumBoundedQueueAdaptor(XeniumBoundedQueueAdaptor&&) = delete; XeniumBoundedQueueAdaptor& operator = (XeniumBoundedQueueAdaptor const&) = delete; XeniumBoundedQueueAdaptor& operator = (XeniumBoundedQueueAdaptor&&) = delete; bool Enqueue(ValueT&& vEnqueue, bool bObligate) { if constexpr(!IsMode(ffpp::QueueFlags::Dynamic)) { auto onBackoff = PolicyT::Backoff(); while(true) { if(uqEnqueued_->try_push(std::move(vEnqueue))) { nSize_.fetch_add(1, std::memory_order::acq_rel); return true; } if(bObligate) { if constexpr(!IsMode(ffpp::QueueFlags::Wait)) break; } else { if constexpr(!IsMode(ffpp::QueueFlags::Wait) || IsMode(ffpp::QueueFlags::Obligate)) break; } onBackoff(); } } else { //Dequeue does not use exclusive lock, while Enqueue requires exclusive access only when internal buffer should grow. SharedLockT sl = this->GetSharedLock(); UniqueLockT ul; while(true) { if(uqEnqueued_->try_push(std::move(vEnqueue))) { nSize_.fetch_add(1, std::memory_order::acq_rel); return true; } size_t nSize = Size(); if(!bObligate && Limit() <= nSize) return false; if(nCapacity_ > nSize) continue; if(sl.owns_lock()) { sl.unlock(); ul = this->GetUniqueLock(); continue; } auto uqIncreased = MakeQueue(nCapacity_ <<= 1); ValueT vPending; while(uqEnqueued_->try_pop(vPending)) if(!uqIncreased->try_push(std::move(vPending))) { return false; } std::swap(uqEnqueued_, uqIncreased); } } return false; } bool Dequeue(ValueT& vDequeued) { SharedLockT sl; //Acquire shared lock only if queue is enabled to grow. if constexpr(IsMode(ffpp::QueueFlags::Dynamic)) sl = this->GetSharedLock(); auto br = uqEnqueued_->try_pop(vDequeued); if(br) nSize_.fetch_sub(1, std::memory_order::acq_rel); return br; } size_t Size() const { return nSize_.load(std::memory_order::acquire); } bool Empty() const { return Size() == 0; } size_t Limit(size_t nLimit) { return nLimit_.exchange(nLimit, std::memory_order::acq_rel); } size_t Limit() { return nLimit_.load(std::memory_order::acquire); } private: using QueueT = TBoundedQueue<ValueT>; using UniqueQueueT = ResourcePolicyT::template UniqueT<QueueT>; static inline size_t const c_nMinimalCapacity_ = 2; UniqueQueueT uqEnqueued_ = MakeQueue(!IsMode(ffpp::QueueFlags::Dynamic) ? FixedCapacity() : c_nMinimalCapacity_); std::atomic<size_t> nLimit_ = c_nInvalidLimit, nSize_ = 0; size_t nCapacity_ = !IsMode(ffpp::QueueFlags::Dynamic) ? FixedCapacity() : c_nMinimalCapacity_; static size_t FixedCapacity() { return std::bit_ceil(std::max(PolicyT::QueueCapacity(), c_nMinimalCapacity_)); } static UniqueQueueT MakeQueue(size_t nCapacity) { return ResourcePolicyT::template AllocateUnique<QueueT>(nCapacity); } };//XeniumBoundedQueueAdaptor #endif//XENIUM_BOUNDED_QUEUE_ADAPTOR_HPP