/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
include/functionalqueue.hpp
861 строка
28 KB
Ant010ff
Version 2.11.3.
08 авг 2026, 16:44
08 авг 2026, 16:44
fd187d3
Код
Авторство
О чём код?
/* This header is part of a Functional Flow Processing Primitives (FFPP) library, version 2.11.3. 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 FFPP_FUNCTIONAL_QUEUE_HPP #define FFPP_FUNCTIONAL_QUEUE_HPP namespace FFPP::Concepts { template<typename Type> concept FunctionalQueuePolicy = FunctionalSequencePolicy<Type> && QueuePolicy<typename Type::QueuePolicyT> && Queue<typename Type::template QueueT<typename Type::ResourcePolicyT::template FunctionT<void()>, typename Type::QueuePolicyT>> && ConstexprInvocableStrict<Type::DoubleQueueing, bool> && InvocableStrict<Type::OnQueueOverflow, bool, uint32_t, Flags> && InvocableStrict<Type::OnQueueCleanupException, bool, std::exception> //Enforcing compile-time use for Type::Queues() static method and result container value_type compatibility: && std::same_as<std::bool_constant<(Type::Queues(), true)>, std::true_type> && QueueIdSet<std::invoke_result_t<decltype(Type::Queues)>> ; }//FFPP::Concepts namespace FFPP { ////FunctionalQueuePolicy///////////////////////////////////////////////////////////////////////////////////////////////////// struct FunctionalQueuePolicy : FunctionalSequencePolicy { using QueuePolicyT = RingQueuePolicy; template<typename TValue, typename TPolicy> using QueueT = RingQueue<TValue, TPolicy>; static constexpr bool DoubleQueueing() { return false; } static constexpr std::vector<uint32_t> Queues() { return { }; } static bool OnQueueOverflow(uint32_t, Flags) { return false; } static bool OnQueueCleanupException(std::exception const&) noexcept { return true; } }; ////FunctionalQueue/////////////////////////////////////////////////////////////////////////////////////////////////////////// template<typename TDerived, Concepts::FunctionalQueuePolicy TPolicy = FunctionalQueuePolicy> class FunctionalQueue : public FunctionalSequence<TDerived, TPolicy> { public: using FunctionalSequenceT = FunctionalSequence<TDerived, TPolicy>; friend FunctionalSequenceT; using ExecutableT = FunctionalSequenceT::ExecutableT; friend ExecutableT; using PolicyT = TPolicy; using ResourcePolicyT = PolicyT::ResourcePolicyT; template<typename Type> using AllocatorT = ResourcePolicyT::template AllocatorT<Type>; using FunctionalT = FunctionalSequenceT::FunctionalT; using LockableT = PolicyT::LockableT; using UniqueLockT = LockableT::UniqueLockT; using SharedLockT = LockableT::SharedLockT; using QueuePolicyT = PolicyT::QueuePolicyT; using QueueT = PolicyT::template QueueT<FunctionalT, QueuePolicyT>; ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// FunctionalQueue( uint32_t nPriority , Flags flgInstance , size_t nLimit = QueueT::c_nInvalidLimit #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) : FunctionalSequenceT( nPriority , flgInstance #if FFPP_TRACK_ORIGIN , slOrigin #endif ) , c_nLimit_(nLimit) , mapSubqueues_([nLimit] () { IdSubqueueMapT mapSubqueues; if constexpr(c_bFixedQueueSet_) { for(auto idQueue : PolicyT::Queues()) { mapSubqueues.emplace( static_cast<uint32_t>(idQueue), ResourcePolicyT::template AllocateUnique<Subqueue>(nLimit) ); } } return mapSubqueues; } ()) { if constexpr(c_bSingleQueue_) { auto itSingle = mapSubqueues_.begin(); psqSingle_ = itSingle->second.get(); idqSingle_ = itSingle->first; } } FunctionalQueue( uint32_t nPriority = PolicyT::DefaultPriority() , size_t nLimit = QueueT::c_nInvalidLimit #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) : FunctionalQueue( nPriority , Flags::Generic , nLimit #if FFPP_TRACK_ORIGIN , slOrigin #endif ) { } FunctionalQueue(FunctionalQueue const&) = delete; FunctionalQueue(FunctionalQueue&&) = delete; FunctionalQueue& operator = (FunctionalQueue const&) = delete; FunctionalQueue& operator = (FunctionalQueue&&) = delete; ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// size_t QueueLimit(uint32_t idQueue, size_t nLimit) { return AcquireQueue(idQueue)->Limit(nLimit); } size_t QueueLimit(uint32_t idQueue) const { auto [pqFront, pqBack] = GetQueues(idQueue); if(pqFront) return pqFront->Limit(); else if(pqBack) return pqBack->Limit(); return 0; } FFPP_ATTR_INLINE size_t Size(uint32_t idQueue = FunctionalSequenceT::InvalidQueueId()) const { if(idQueue == FunctionalSequenceT::InvalidQueueId()) { //nEnqueued_ can transiently underflow. This is not harmful, but may temporally influence on thread pool dispatch //and load balancing efficiency (FunctionalQueuePol). using CounterT = FunctionalSequenceT::AtomicWaitableCounterT::value_type; CounterT nEnqueued = this->nEnqueued_.load(std::memory_order::relaxed); CounterT constexpr c_mskMSB = ~(std::numeric_limits<CounterT>::max() >> 1); return nEnqueued < c_mskMSB ? nEnqueued : 0; //expecting conditional move instruction here } else { auto [pqFront, pqBack] = GetQueues(idQueue); size_t nSize = 0; if(pqFront) nSize += pqFront->Size(); if(pqBack) nSize += pqBack->Size(); return nSize; } } size_t Count() const { size_t nCount = 0; if constexpr(c_bFixedQueueSet_) { if constexpr(c_bSingleQueue_) { auto [pqFront, pqBack] = psqSingle_->GetQueues(); nCount += pqFront->Size(); if constexpr(c_bDoubleQueueing_) nCount += pqBack->Size(); } else { for(auto const& [_, upSubqueue] : mapSubqueues_) { auto [pqFront, pqBack] = upSubqueue->GetQueues(); nCount += pqFront->Size(); if constexpr(c_bDoubleQueueing_) nCount += pqBack->Size(); } } } else { auto sl = this->GetSharedLock(); for(auto const& [_, upSubqueue] : mapSubqueues_) { auto [pqFront, pqBack] = upSubqueue->GetQueues(); nCount += pqFront->Size(); if constexpr(c_bDoubleQueueing_) nCount += pqBack->Size(); } } return nCount; } size_t Limit() const noexcept { return c_nLimit_; } protected: using FunctionalSequenceT::OnEnqueued; using FunctionalSequenceT::OnDequeued; using FunctionalSequenceT::OnSubmit; using FunctionalSequenceT::OnWait; using FunctionalSequenceT::HasEnqueued; using FunctionalSequenceT::WaitEnqueue; template<std::invocable<TDerived*> TFunctional> FFPP_ATTR_HOT_PATH bool OnEnqueue(TFunctional&& fnTask, uint32_t idQueue, Flags flgContext) { bool const bObligate = !!(flgContext & Flags::Obligate); bool const bEnqueued = AcquireFrontQueue(idQueue)->Enqueue(std::move(fnTask), bObligate); if(!bEnqueued && !PolicyT::OnQueueOverflow(idQueue, flgContext)) { Except<std::overflow_error>::Throw( "Queue (", idQueue, ") overflow,", bObligate ? "obligate @" : "regular @" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); } return bEnqueued; } template<std::invocable<TDerived*> TFunctional> std::pair<bool, uint32_t> OnComplete(Flags flgContext, TFunctional&& onComplete) { auto [idQueue, pQueue] = AcquireCompletionQueue(!!(flgContext & Flags::Priority)); if(!pQueue) Except<>::Throw( "Failed to acquire queue @" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); bool bEnqueued = pQueue->Enqueue(std::forward<TFunctional>(onComplete), true); if(!bEnqueued && !PolicyT::OnQueueOverflow(idQueue, flgContext | Flags::Obligate)) { Except<std::overflow_error>::Throw( "Queue (", idQueue, ") overflow on completion attempt @" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); } return { bEnqueued, idQueue }; } FFPP_ATTR_HOT_PATH size_t OnProcess(bool bAll, uint32_t idQueue) { size_t nProcessed = 0; if(idQueue == FunctionalSequenceT::InvalidQueueId()) { if(bAll) nProcessed = OnProcessAll(); else nProcessed = OnProcessSingle(); } else { if(bAll) nProcessed = OnProcessAll(idQueue); else nProcessed = OnProcessSingle(idQueue); } return nProcessed; } void OnFinalize() noexcept { Clear(); FunctionalSequenceT::OnFinalize(); } std::pair<size_t, bool> Clear(std::convertible_to<uint32_t> auto... idQueues) { using TDerivedPointer = TDerived*; CompatibleCounterT nClear = 0; bool bException = false; FFPP_TRY { //Regular cleanups are possible only in context of instance been cleared (usually Actor or FunctionalQueueThread) //or once from any other thread (in case of concurrent Finalize method calls). if(!(this->IsProcessingThread() || bClearOnce_.exchange(false, std::memory_order::acq_rel))) { Except<>::Throw( "Invalid cleanup attempt @" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); } std::vector<QueueT*, AllocatorT<QueueT*>> vClear; bool constexpr c_bLock = !c_bFixedQueueSet_; SharedLockT sl; if constexpr(c_bLock) sl = this->GetSharedLock(); if constexpr(sizeof...(idQueues) > 0) { if constexpr(c_bDoubleQueueing_) vClear.reserve(sizeof...(idQueues) * 2); else vClear.reserve(sizeof...(idQueues)); for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) if(((idSubqueue == idQueues) || ...)) { auto [pqFront, pqBack] = upSubqueue->GetQueues(); vClear.push_back(pqFront); if constexpr(c_bDoubleQueueing_) vClear.push_back(pqBack); } } else { if constexpr(c_bDoubleQueueing_) vClear.reserve(mapSubqueues_.size() * 2); else vClear.reserve(mapSubqueues_.size()); for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) { auto [pqFront, pqBack] = upSubqueue->GetQueues(); vClear.push_back(pqFront); if constexpr(c_bDoubleQueueing_) vClear.push_back(pqBack); } } if constexpr(c_bLock) sl.unlock(); FunctionalT fnPending; for(auto* pQueue : vClear) if(pQueue) { while(pQueue->Dequeue(fnPending)) { nClear++; if constexpr(c_bQueueSwapSemantics_) fnPending = { }; } } } FFPP_CATCH(std::exception, ex) { if constexpr(requires { PolicyT::OnQueueCleanupException(TDerivedPointer { }, std::exception { }); }) { PolicyT::OnQueueCleanupException(static_cast<TDerived*>(this), ex); } else { PolicyT::OnQueueCleanupException(ex); } bException = true; } if(nClear > 0) { FFPP_TRY { static_cast<TDerived*>(this)->template OnDequeued<false>(FunctionalSequenceT::InvalidQueueId(), nClear); } FFPP_CATCH(std::exception, ex) { if constexpr(requires { PolicyT::OnQueueCleanupException(TDerivedPointer { }, std::exception { }); }) { PolicyT::OnQueueCleanupException(static_cast<TDerived*>(this), ex); } else { PolicyT::OnQueueCleanupException(ex); } bException = true; } } return { nClear, bException }; } FFPP_ATTR_HOT_PATH uint32_t Steal(FunctionalT& fnTask) { auto sgAccessCounter = this->BeginTrackContention(); if constexpr(c_bDoubleQueueing_) { //In case of double-queueing try to dequeue (steal) from current front queue segment (pqFront) which should not be //contended by processing thread. If optimistic attempt failed try to dequeue from back queue segment (pqBack). auto [idqPending, pqFront, pqBack] = GetFrontBackQueues(); if((pqFront && pqFront->Dequeue(fnTask)) || (pqBack && pqBack->Dequeue(fnTask))) { static_cast<TDerived*>(this)->template OnDequeued<false>(idqPending); return idqPending; } } else { auto [idqPending, pqPending] = GetBackQueue(); if(!pqPending) return FunctionalSequenceT::InvalidQueueId(); if(pqPending->Dequeue(fnTask)) { static_cast<TDerived*>(this)->template OnDequeued<false>(idqPending); return idqPending; } } return FunctionalSequenceT::InvalidQueueId(); } FFPP_ATTR_INLINE size_t OnProcessAll(uint32_t idQueue) { UniqueQueueT upQueue; QueueT* pQueue = nullptr; if constexpr(c_bFixedQueueSet_) { if constexpr(c_bSingleQueue_) { if(idQueue != idqSingle_) [[unlikely]] return 0; pQueue = psqSingle_->GetBackQueue(); } else { auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() == itSubqueue) return 0; pQueue = itSubqueue->second->GetBackQueue(); } } else { auto ul = this->GetUniqueLock(); auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() == itSubqueue) return 0; upQueue = itSubqueue->second->MoveBackQueue(); pQueue = upQueue.get(); mapSubqueues_.erase(itSubqueue); } auto const nSize = pQueue->Size(); if(nSize == 0) return 0; auto* pThis = static_cast<TDerived*>(this); CompatibleCounterT nProcessed = 0; UniqueGuard ug { [&nProcessed, idQueue, pThis] () noexcept { if(nProcessed > 0) pThis->template OnDequeued<false>(idQueue, nProcessed); } }; FunctionalT fnPending; while(pQueue->Dequeue(fnPending) && nProcessed < nSize) { std::invoke(std::move(fnPending), pThis); nProcessed++; if constexpr(c_bQueueSwapSemantics_) fnPending = { }; } ug.Abandon(); if(nProcessed > 0) pThis->template OnDequeued<true>(idQueue, nProcessed); return nProcessed; } FFPP_ATTR_INLINE size_t OnProcessAll() { IdSubqueueMapT mapSubqueues; IdSubqueueMapCondT* pmapSubqueues = nullptr; if constexpr(c_bFixedQueueSet_) { pmapSubqueues = &mapSubqueues_; } else { auto ul = this->GetUniqueLock(); mapSubqueues = std::move(mapSubqueues_); pmapSubqueues = &mapSubqueues; } CompatibleCounterT nProcessed = 0; auto* pThis = static_cast<TDerived*>(this); UniqueGuard ug { [&nProcessed, pThis] () noexcept { if(nProcessed > 0) pThis->template OnDequeued<false>(FunctionalSequenceT::InvalidQueueId(), nProcessed); } }; FunctionalT fnPending; for(auto const& [idSubqueue, upSubqueue] : *pmapSubqueues) { QueueT* pQueue = upSubqueue->GetBackQueue(); size_t nSize = pQueue->Size(); if(nSize == 0) continue; while(pQueue->Dequeue(fnPending) && nSize-- > 0) { std::invoke(std::move(fnPending), pThis); nProcessed++; if constexpr(c_bQueueSwapSemantics_) fnPending = { }; } } ug.Abandon(); if(nProcessed > 0) pThis->template OnDequeued<true>(FunctionalSequenceT::InvalidQueueId(), nProcessed); return nProcessed; } FFPP_ATTR_INLINE size_t OnProcessSingle(uint32_t idQueue) { auto pQueue = GetBackQueue(idQueue); return OnProcessSingle(idQueue, pQueue); } FFPP_ATTR_INLINE size_t OnProcessSingle() { auto [idQueue, pQueue] = GetBackQueue(); return OnProcessSingle(idQueue, pQueue); } private: static constexpr bool c_bDoubleQueueing_ = PolicyT::DoubleQueueing() , c_bSingleQueue_ = PolicyT::Queues().size() == 1 , c_bFixedQueueSet_ = PolicyT::Queues().size() > 0 , c_bConstSubqueueMap = c_bFixedQueueSet_ && !c_bDoubleQueueing_ , c_bQueueSwapSemantics_ = [] () consteval { if constexpr(requires { { QueueT::IsSwapSemantics() } -> std::same_as<bool>; }) { return QueueT::IsSwapSemantics(); } return false; } () ; static constexpr uint32_t c_idPriorityQueue_ = std::numeric_limits<uint32_t>::max() - 1; using UniqueQueueT = typename ResourcePolicyT::template UniqueT<QueueT>; struct Subqueue { static constexpr size_t c_mskQueue = 1; using DoubleQueueArrayT = std::array<UniqueQueueT, 2>; DoubleQueueArrayT const c_arDoubleQueue; std::atomic<size_t> nQueue = 0; Subqueue(size_t nQueueLimit) : c_arDoubleQueue([nQueueLimit] () -> DoubleQueueArrayT { if constexpr(c_bDoubleQueueing_) { return { ResourcePolicyT::template AllocateUnique<QueueT>(nQueueLimit), ResourcePolicyT::template AllocateUnique<QueueT>(nQueueLimit) }; } else { return { ResourcePolicyT::template AllocateUnique<QueueT>(nQueueLimit), UniqueQueueT { } }; } } ()) { } FFPP_ATTR_INLINE void SwapQueues() noexcept { static_assert(PolicyT::DoubleQueueing(), "Must be used only in a double-queueing context!"); nQueue.fetch_add(1, std::memory_order::release); } FFPP_ATTR_INLINE std::pair<QueueT*, QueueT*> GetQueues() const noexcept { return { c_arDoubleQueue[0].get(), c_arDoubleQueue[1].get() }; } FFPP_ATTR_INLINE std::pair<QueueT*, QueueT*> GetFrontBackQueues() const noexcept { static_assert(PolicyT::DoubleQueueing(), "Must be used only in a double-queueing context!"); size_t const nCurrent = nQueue.load(std::memory_order::acquire); return { c_arDoubleQueue[nCurrent & c_mskQueue].get(), c_arDoubleQueue[(nCurrent + 1) & c_mskQueue].get() }; } FFPP_ATTR_INLINE QueueT* GetFrontQueue() const noexcept { if constexpr(c_bDoubleQueueing_) { size_t const iQueue = nQueue.load(std::memory_order::acquire) & c_mskQueue; return c_arDoubleQueue[iQueue].get(); } else { return c_arDoubleQueue[0].get(); } } FFPP_ATTR_INLINE QueueT* GetBackQueue() const noexcept { if constexpr(c_bDoubleQueueing_) { size_t const iQueue = (nQueue.load(std::memory_order::acquire) + 1) & c_mskQueue; return c_arDoubleQueue[iQueue].get(); } else { return c_arDoubleQueue[0].get(); } } FFPP_ATTR_INLINE UniqueQueueT MoveBackQueue() noexcept { static_assert(PolicyT::Queues().size() == 0, "Must be used only in a dynamic queue-set context!"); if constexpr(c_bDoubleQueueing_) { size_t const iQueue = (nQueue.load(std::memory_order::acquire) + 1) & c_mskQueue; return std::move(c_arDoubleQueue[iQueue]); } else { return std::move(c_arDoubleQueue[0]); } } };//Subqueue using UniqueSubqueueT = typename ResourcePolicyT::template UniqueT<Subqueue>; #if defined(__cpp_lib_flat_map) && __cpp_lib_flat_map >= 202207L using IdSubqueueMapT = std::flat_map< uint32_t, UniqueSubqueueT, std::greater<uint32_t>, std::vector<uint32_t, typename ResourcePolicyT::template AllocatorT<uint32_t>>, std::vector<UniqueSubqueueT, typename ResourcePolicyT::template AllocatorT<UniqueSubqueueT>> >; #elif defined(BOOST_CONTAINER_FLAT_MAP_HPP) using IdSubqueuePairT = std::pair<uint32_t, UniqueSubqueueT>; using IdSubqueueMapT = boost::container::flat_map< uint32_t, UniqueSubqueueT, std::greater<uint32_t>, std::vector<IdSubqueuePairT, typename ResourcePolicyT::template AllocatorT<IdSubqueuePairT>> >; #else using IdSubqueuePairT = std::pair<uint32_t const, UniqueSubqueueT>; using IdSubqueueMapT = std::map< uint32_t, UniqueSubqueueT, std::greater<uint32_t>, typename ResourcePolicyT::template AllocatorT<IdSubqueuePairT> >; #endif using IdSubqueueMapCondT = std::conditional_t<c_bConstSubqueueMap, IdSubqueueMapT const, IdSubqueueMapT>; size_t const c_nLimit_ = QueueT::c_nInvalidLimit; std::atomic<bool> bClearOnce_ = true; IdSubqueueMapCondT mapSubqueues_; uint32_t idqSingle_ = FunctionalSequenceT::InvalidQueueId(); Subqueue* psqSingle_ = nullptr; FFPP_ATTR_INLINE std::pair<QueueT*, QueueT*> GetQueues(uint32_t idQueue) const { SharedLockT sl; if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); if constexpr(c_bSingleQueue_) { if(idQueue == idqSingle_) [[likely]] return psqSingle_->GetQueues(); } else { auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() != itSubqueue) return itSubqueue->second->GetQueues(); } return { nullptr, nullptr }; } FFPP_ATTR_INLINE QueueT* AcquireQueue(uint32_t idQueue) { if constexpr(c_bFixedQueueSet_) { if constexpr(c_bSingleQueue_) { if(idqSingle_ == idQueue) [[likely]] return psqSingle_->GetFrontQueue(); } else { auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() != itSubqueue) return itSubqueue->second->GetFrontQueue(); } } else { UniqueLockT ul; SharedLockT sl = this->GetSharedLock(); do { auto itSubqueue = mapSubqueues_.lower_bound(idQueue); if(mapSubqueues_.end() == itSubqueue || itSubqueue->first != idQueue) { if(sl.owns_lock()) { sl.unlock(); ul = this->GetUniqueLock(); continue; } itSubqueue = mapSubqueues_.emplace_hint( itSubqueue, idQueue, ResourcePolicyT::template AllocateUnique<Subqueue>(c_nLimit_) ); } return itSubqueue->second->GetFrontQueue(); } while(true); } Except<>::Throw( "Failed to acquire queue", idQueue, "@" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); return nullptr; } std::pair<uint32_t, QueueT*> AcquireCompletionQueue(bool bPriority) { if constexpr(c_bFixedQueueSet_) { if constexpr(c_bSingleQueue_) { return { idqSingle_, psqSingle_->GetFrontQueue() }; } else { if(bPriority) { auto const& [idSubqueue, upSubqueue] = *mapSubqueues_.begin(); return { idSubqueue, upSubqueue->GetFrontQueue() }; } else { auto const& [idSubqueue, upSubqueue] = *mapSubqueues_.rbegin(); return { idSubqueue, upSubqueue->GetFrontQueue() }; } } } else { auto sl = this->GetSharedLock(); if(bPriority || mapSubqueues_.empty()) { sl.unlock(); return { c_idPriorityQueue_, AcquireQueue(c_idPriorityQueue_) }; } else { auto const& [idSubqueue, upSubqueue] = *mapSubqueues_.rbegin(); return { idSubqueue, upSubqueue->GetFrontQueue() }; } } } FFPP_ATTR_INLINE QueueT* AcquireFrontQueue(uint32_t idQueue) { if constexpr(c_bFixedQueueSet_) { if constexpr(c_bSingleQueue_) { if(idQueue == idqSingle_) [[likely]] return psqSingle_->GetFrontQueue(); } else { auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() != itSubqueue) return itSubqueue->second->GetFrontQueue(); } } else { UniqueLockT ul; SharedLockT sl = this->GetSharedLock(); do { auto itSubqueue = mapSubqueues_.lower_bound(idQueue); if(mapSubqueues_.end() == itSubqueue || itSubqueue->first != idQueue) { if(sl.owns_lock()) { sl.unlock(); ul = this->GetUniqueLock(); continue; } itSubqueue = mapSubqueues_.emplace_hint( itSubqueue, idQueue, ResourcePolicyT::template AllocateUnique<Subqueue>(c_nLimit_) ); } return itSubqueue->second->GetFrontQueue(); } while(true); } Except<>::Throw( "Failed to acquire front queue", idQueue, "@" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); return nullptr; } FFPP_ATTR_INLINE std::pair<uint32_t, QueueT*> GetBackQueue() { //It is safe to return potentially empty queue even in case of multiple consumers (DispatchFlags::Steal), some of them //will do an idle processing-loop iteration. Generally, idle processing-loop iterations considered not to be an issue. if constexpr(c_bDoubleQueueing_) { if constexpr(c_bSingleQueue_) { //Double-queueing with a single queue acts not only like a contention mitigator but also a bulk-processing //facility. auto [pqFront, pqBack] = psqSingle_->GetFrontBackQueues(); if(pqBack->Empty()) { if(!pqFront->Empty()) { psqSingle_->SwapQueues(); return { idqSingle_, pqFront }; } } else return { idqSingle_, pqBack }; return { FunctionalSequenceT::InvalidQueueId(), nullptr }; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); //Here two strategies are possible: //- iterate mapSubqueues_ every time and return currently non-empty back-queue; //- cache last non-empty subqueue in a member pointer and update it only when its back-queue becomes empty, but //this may require exclusive lock when work-stealing enabled (DispatchFlags::Steal). Considering that mapSubqueues_ //usually contains only few entries and an accent on task priorities the first strategy was chosen. for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) { auto [pqFront, pqBack] = upSubqueue->GetFrontBackQueues(); if(!pqBack->Empty()) { return { idSubqueue, pqBack }; } else if(!pqFront->Empty()) { upSubqueue->SwapQueues(); return { idSubqueue, pqFront }; } } } else { if constexpr(c_bSingleQueue_) { QueueT* pqBack = psqSingle_->GetBackQueue(); if(!pqBack->Empty()) return { idqSingle_, pqBack }; return { FunctionalSequenceT::InvalidQueueId(), nullptr }; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) { QueueT* pQueue = upSubqueue->GetBackQueue(); if(!pQueue->Empty()) return { idSubqueue, pQueue }; } } return { FunctionalSequenceT::InvalidQueueId(), nullptr }; } FFPP_ATTR_INLINE QueueT* GetBackQueue(uint32_t idQueue) { if constexpr(c_bDoubleQueueing_) { if constexpr(c_bSingleQueue_) { if(idqSingle_ != idQueue) return nullptr; auto [pqFront, pqBack] = psqSingle_->GetFrontBackQueues(); if(pqBack->Empty()) { if(!pqFront->Empty()) { psqSingle_->SwapQueues(); return pqFront; } } else return pqBack; return nullptr; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() == itSubqueue) return nullptr; auto [pqFront, pqBack] = itSubqueue->second->GetFrontBackQueues(); if(!pqBack->Empty()) { return pqBack; } else if(!pqFront->Empty()) { itSubqueue->second->SwapQueues(); return pqFront; } } else { if constexpr(c_bSingleQueue_) { if(idqSingle_ != idQueue) [[unlikely]] return nullptr; QueueT* pqBack = psqSingle_->GetBackQueue(); if(!pqBack->Empty()) return pqBack; return nullptr; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); auto itSubqueue = mapSubqueues_.find(idQueue); if(mapSubqueues_.end() == itSubqueue) return nullptr; QueueT* pQueue = itSubqueue->second->GetBackQueue(); if(!pQueue->Empty()) return pQueue; } return nullptr; } FFPP_ATTR_INLINE std::tuple<uint32_t, QueueT*, QueueT*> GetFrontBackQueues() { if constexpr(c_bDoubleQueueing_) { if constexpr(c_bSingleQueue_) { auto [pqFront, pqBack] = psqSingle_->GetFrontBackQueues(); if(!pqFront->Empty() || !pqBack->Empty()) return { idqSingle_, pqFront, pqBack }; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) { auto [pqFront, pqBack] = upSubqueue->GetFrontBackQueues(); if(!pqFront->Empty() || !pqBack->Empty()) return { idSubqueue, pqFront, pqBack }; } } else { if constexpr(c_bSingleQueue_) { QueueT* pQueue = psqSingle_->GetBackQueue(); if(!pQueue->Empty()) return { idqSingle_, pQueue, pQueue }; return { FunctionalSequenceT::InvalidQueueId(), nullptr, nullptr }; } SharedLockT sl; //!c_bFixedQueueSet_ if constexpr(!c_bFixedQueueSet_) sl = this->GetSharedLock(); for(auto const& [idSubqueue, upSubqueue] : mapSubqueues_) { QueueT* pQueue = upSubqueue->GetBackQueue(); if(!pQueue->Empty()) return { idqSingle_, pQueue, pQueue }; } } return { FunctionalSequenceT::InvalidQueueId(), nullptr, nullptr }; } FFPP_ATTR_INLINE size_t OnProcessSingle(uint32_t idQueue, QueueT* pQueue) { if(pQueue == nullptr) return 0; auto* pThis = static_cast<TDerived*>(this); bool bDequeued = false; UniqueGuard ug { [&bDequeued, idQueue, pThis] () noexcept { if(bDequeued) pThis->template OnDequeued<false>(idQueue); } }; FunctionalT fnPending; bDequeued = pQueue->Dequeue(fnPending); if(!bDequeued) [[unlikely]] { ug.Abandon(); return 0; } std::invoke(fnPending, pThis); ug.Abandon(); pThis->template OnDequeued<true>(idQueue); return 1; } };//FunctionalQueue }//FFPP #endif//FFPP_FUNCTIONAL_QUEUE_HPP