/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
include/functionalloop.hpp
208 строк
6 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_LOOP_HPP #define FFPP_FUNCTIONAL_LOOP_HPP namespace FFPP::Concepts { template<typename Type> concept FunctionalLoopPolicy = FunctionalSequencePolicy<Type> && SequencePolicy<Type>; }//FFPP::Concepts namespace FFPP { ////FunctionalLoopPolicy////////////////////////////////////////////////////////////////////////////////////////////////////// struct FunctionalLoopPolicy : FunctionalSequencePolicy, SequencePolicy { static constexpr SchedulingFlags SchedulingMethod() { return SchedulingFlags::Default; } }; ////FunctionalLoop//////////////////////////////////////////////////////////////////////////////////////////////////////////// template<typename TDerived, Concepts::FunctionalLoopPolicy TPolicy = FunctionalLoopPolicy> class FunctionalLoop : 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; using FunctionalT = FunctionalSequenceT::FunctionalT; ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// FunctionalLoop( uint32_t nPriority = 0 , Flags flgSequence = Flags::Generic #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) : FunctionalSequenceT( nPriority , flgSequence #if FFPP_TRACK_ORIGIN , slOrigin #endif ) , bClosed_(!!(flgSequence & Flags::Commit)) , bCommitted_(!!(flgSequence & Flags::Commit)) { } FunctionalLoop(FunctionalLoop const&) = delete; FunctionalLoop(FunctionalLoop&&) = delete; FunctionalLoop& operator = (FunctionalLoop const&) = delete; FunctionalLoop& operator = (FunctionalLoop&&) = delete; ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// bool Commit() noexcept { bool bClosed = bClosed_.exchange(true, std::memory_order::relaxed); if(!bClosed) { assert(Size() > 0); fmodSize_ = { Size() }; bCommitted_.store(true, std::memory_order::release); } return bClosed; } [[nodiscard]] bool IsCommitted() const noexcept { return bCommitted_.load(std::memory_order::acquire); } [[nodiscard]] size_t Size() const noexcept { return tsLoop_.Size(); } [[nodiscard]] size_t Progress() const noexcept { return iLoop_.load(std::memory_order::relaxed); } [[nodiscard]] size_t Loop() const noexcept { return nLoop_.load(std::memory_order::relaxed); } [[nodiscard]] bool Ready() const noexcept { return (IsCommitted() && Size() > 0) ? fmodSize_(iLoop_.load(std::memory_order::relaxed)) == 0 : false; } [[nodiscard]] bool Begin() const noexcept { return (IsCommitted() && Size() > 0) ? fmodSize_(iLoop_.load(std::memory_order::relaxed)) == 1 : false; } void Continue() noexcept { size_t nSize = IsCommitted() ? Size() : 0; if(nSize > 0) { size_t iLoop = iLoop_.load(std::memory_order::relaxed); while(!iLoop_.compare_exchange_weak( iLoop, (iLoop + nSize) - fmodSize_(iLoop), std::memory_order::relaxed, std::memory_order::relaxed )); } } protected: using FunctionalSequenceT::OnDequeued; using FunctionalSequenceT::HasEnqueued; using FunctionalSequenceT::WaitEnqueue; bool OnSubmit() { return IsCommitted() && FunctionalSequenceT::OnSubmit(); } template<bool t_bNotify = false> //FunctionalLoop's enqueue counter is not used for any thread notification. std::pair<bool, size_t> OnEnqueued(uint32_t idQueue) { return FunctionalSequenceT::template OnEnqueued<t_bNotify>(idQueue); } template<std::invocable<TDerived*> TFunctional> bool OnEnqueue(TFunctional&& fnTask, uint32_t, Flags) { //Actually, enqueues into FunctionalLoop supposed to be synchronous in a single thread constructing the loop (or //synchronized externally, not required in FFPP itself). The guard is harmless and can be helpful during debugging. if(bClosed_.load(std::memory_order::relaxed)) { assert(false); return false; } return tsLoop_.Append(std::move(fnTask)); } template<std::invocable<TDerived*> TFunctional> std::pair<bool, uint32_t> OnComplete([[maybe_unused]] Flags flgContext, TFunctional&& onComplete) { return { tsLoop_.Append(std::forward<TFunctional>(onComplete)), 0 }; } size_t OnProcess(bool bAll, uint32_t) { this->UpdateProcessingThread(); size_t nSize = Size(), nProcessed = 0; if(nSize > 0) { if(bAll) nProcessed = OnProcessAll(nSize); //Sequential invocation supposed, no concurrency here. else nProcessed = OnProcessSingle(nSize); } return nProcessed; } size_t OnProcessAll() { this->UpdateProcessingThread(); return OnProcessAll(Size()); } size_t OnProcessSingle() { //Sequential invocation supposed, no concurrency here. this->UpdateProcessingThread(); return OnProcessSingle(Size()); } void OnFinalize() noexcept { FunctionalSequenceT::OnFinalize(); } FFPP_ATTR_INLINE size_t OnProcessAll([[maybe_unused]] size_t nSize) { iLoop_.fetch_add(nSize, std::memory_order::relaxed); nLoop_.fetch_add(1, std::memory_order::relaxed); tsLoop_.Visit( [] (auto& fnTask, auto* pThis) { fnTask(pThis); }, static_cast<TDerived*>(this) ); return nSize; } FFPP_ATTR_INLINE size_t OnProcessSingle([[maybe_unused]] size_t nSize) { size_t const iLoop = fmodSize_(iLoop_.fetch_add(1, std::memory_order::relaxed)); if(iLoop == 0) nLoop_.fetch_add(1, std::memory_order::relaxed); tsLoop_.Visit( iLoop, [] (auto& fnTask, auto* pThis) { fnTask(pThis); }, static_cast<TDerived*>(this) ); return 1; } private: using TaskSequenceT = Sequence<FunctionalT, PolicyT>; TaskSequenceT tsLoop_; math::FastMod fmodSize_; std::atomic<size_t> iLoop_ = 0, //loop iteration index nLoop_ = 0; //current loop iteration, two counters simplifies overflow handling (see dtl::BrokerTaskLoop, BrokerChainer) std::atomic<bool> bClosed_ = false, bCommitted_ = false; };//FunctionalLoop }//FFPP #endif//FFPP_FUNCTIONAL_LOOP_HPP