/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
include/functionalsequence.hpp
382 строки
14 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_SEQUENCE_HPP #define FFPP_FUNCTIONAL_SEQUENCE_HPP namespace FFPP::Concepts { template<typename Type> concept FunctionalSequencePolicy = ExecutablePolicy<Type> && ConstexprInvocableStrict<Type::TrackContention, bool> //Enforcing compile-time use for Type::DefaultQueue() static method and result type compatibility: && std::same_as<std::bool_constant<(Type::DefaultQueue(), true)>, std::true_type> && QueueId<typename std::invoke_result_t<decltype(Type::DefaultQueue)>> ; template<typename Type> concept FunctionalSequence = FunctionalSequencePolicy<typename Type::PolicyT> && Executable<Type> && requires { typename Type::FunctionalT; } && requires(Type* pFS) { { pFS->GetSequenceFlags() } -> std::same_as<Flags>; { pFS->IsSequence(Flags(), bool()) } -> std::same_as<bool>; { pFS->Enqueue(typename Type::FunctionalT(), uint32_t(), Flags { }) } -> std::same_as<bool>; } ; }//FFPP::Concepts namespace FFPP { ////FunctionalSequencePolicy////////////////////////////////////////////////////////////////////////////////////////////////// struct FunctionalSequencePolicy : ExecutablePolicy { template<bool t_bAutosubmit, typename TArguments, typename TFunctionalSequence> using SequenceChainerT = SequenceChainer<t_bAutosubmit, TArguments, TFunctionalSequence>; static constexpr bool TrackContention() { return false; } static constexpr uint32_t DefaultQueue() { return 0; } }; ////FunctionalSequence//////////////////////////////////////////////////////////////////////////////////////////////////////// template<typename TDerived, Concepts::FunctionalSequencePolicy TPolicy = FunctionalSequencePolicy> class FunctionalSequence : public Executable<TDerived, TPolicy> { public: using PolicyT = TPolicy; using ResourcePolicyT = PolicyT::ResourcePolicyT; template<typename TSignature, bool t_bUnique = true> using FunctionT = ResourcePolicyT::template FunctionT<TSignature, t_bUnique>; using FunctionalT = FunctionT<void(TDerived*)>; template<typename Type> using SharedT = ResourcePolicyT::template SharedT<Type>; template<typename Type> using WeakT = ResourcePolicyT::template WeakT<Type>; using SP = SharedT<TDerived>; using WP = WeakT<TDerived>; template<bool t_bAutosubmit, Concepts::Tuple TArguments = std::tuple<>> using SequenceChainerT = PolicyT::template SequenceChainerT<t_bAutosubmit, TArguments, TDerived>; using ExecutableT = Executable<TDerived, TPolicy>; static_assert( !TestModeFlags(PolicyT::SchedulingMethod(), SchedulingFlags::Concurrent, false) , "Concurrent scheduling modes are not applicable to functional sequence objects." ); ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// FunctionalSequence( uint32_t nPriority = PolicyT::DefaultPriority() , Flags flgSequence = Flags::Generic #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) : ExecutableT( nPriority , flgSequence #if FFPP_TRACK_ORIGIN , slOrigin #endif ) { } FunctionalSequence(FunctionalSequence const&) = delete; FunctionalSequence(FunctionalSequence&&) = delete; FunctionalSequence& operator = (FunctionalSequence const&) = delete; FunctionalSequence& operator = (FunctionalSequence&&) = delete; ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// FFPP_ATTR_INLINE static SharedT<TDerived> MakeShared( uint32_t nPriority = PolicyT::DefaultPriority() , Flags flgSequence = Flags::Generic #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) { return ResourcePolicyT::template AllocateShared<TDerived>( nPriority , flgSequence #if FFPP_TRACK_ORIGIN , slOrigin #endif ); } FFPP_ATTR_INLINE Flags GetSequenceFlags() const noexcept { return this->GetInstanceFlags(); } FFPP_ATTR_INLINE bool IsSequence(Flags flg, bool bAll = false) const noexcept { return this->TestInstanceFlags(flg, bAll); } template<bool t_bInternalInvocation = false> Flags AddSequenceFlags(Flags flgAdd) noexcept { return this->template AddInstanceFlags<t_bInternalInvocation>(flgAdd); } template<bool t_bInternalInvocation = false> Flags ClearSequenceFlags(Flags flgClear) noexcept { return this->template ClearInstanceFlags<t_bInternalInvocation>(flgClear); } ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// template<std::invocable<TDerived*> TFunctional> FFPP_ATTR_HOT_PATH bool Enqueue(TFunctional&& fnTask, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { if(this->IsComplete()) return false; CounterGuardT sgAccessCounter = BeginTrackContention(); if(!static_cast<TDerived*>(this)->OnEnqueue(std::move(fnTask), idQueue, flgContext)) return false; return std::get<bool>(static_cast<TDerived*>(this)->OnEnqueued(idQueue)); } template<typename TMessage> bool Enqueue(TMessage&& msg, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return Enqueue( [msg = std::move(msg)] (TDerived* pSelf) mutable { if constexpr(requires { pSelf->On(std::move(msg)); }) pSelf->On(std::move(msg)); else if constexpr(requires { pSelf->On(msg); }) pSelf->On(msg); else static_assert(requires { pSelf->On(msg); }, "Undefined message handler!"); }, idQueue, flgContext ); } template<typename TMessage> bool Enqueue(TMessage const& msg, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return Enqueue( [msg = msg] (TDerived* pSelf) mutable { if constexpr(requires { pSelf->On(std::move(msg)); }) pSelf->On(std::move(msg)); else if constexpr(requires { pSelf->On(msg); }) pSelf->On(msg); else static_assert(requires { pSelf->On(msg); }, "Undefined message handler!"); }, idQueue, flgContext ); } ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// template<std::invocable<TDerived*> TFunctional> bool Complete(Flags flgContext, TFunctional&& onComplete) { bool bComplete = this->ExecutableT::TryComplete(); if(bComplete) { TDerived* pSelf = static_cast<TDerived*>(this); UniqueGuard<> ug { [pSelf, flgContext] () { pSelf->Finalize(flgContext & ~Flags::Wait); } }; FunctionalT fnComplete = std::move(onComplete); auto [bEnqueued, idQueue] = !!fnComplete ? pSelf->OnComplete( flgContext, [fnComplete = std::move(fnComplete), flgContext] (TDerived* pSelf) mutable { UniqueGuard<> ug { [pSelf, flgContext] () { pSelf->Finalize(flgContext & ~Flags::Wait); } }; std::invoke(std::move(fnComplete), pSelf); } ) : pSelf->OnComplete( flgContext, [flgContext] (TDerived* pSelf) { pSelf->Finalize(flgContext & ~Flags::Wait); } ) ; if(bEnqueued) { bEnqueued = pSelf->OnEnqueued(idQueue).first; } if(bEnqueued) { ug.Abandon(); } else { flgContext = flgContext & ~Flags::Wait; bComplete = false; } } if(!!(flgContext & Flags::Wait)) this->Wait(flgContext); return bComplete; } template<std::invocable<TDerived*> TFunctional> bool Complete(TFunctional&& onComplete) { return Complete(Flags::Undefined, std::forward<TFunctional>(onComplete)); } bool Complete(Flags flgContext = Flags::Undefined) { return Complete(flgContext, FunctionalT { }); } ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// template<std::invocable<TDerived*> TFunctional> FFPP_ATTR_HOT_PATH auto Emit(TFunctional&& fnTask, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return SequenceChainerT<false> { this->Self(), !(flgContext & Flags::Deferred), Enqueue(std::forward<TFunctional>(fnTask), idQueue, flgContext) }; } template<typename TMessage> auto Emit(TMessage&& msg, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return SequenceChainerT<false> { this->Self(), !(flgContext & Flags::Deferred), Enqueue(std::forward<TMessage>(msg), idQueue, flgContext) }; } template<std::invocable TFunctional> auto Chain(TFunctional&& fnTask, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return SequenceChainerT<true> { this->Self(), !(flgContext & Flags::Deferred), true } .Chain(std::forward<TFunctional>(fnTask), idQueue, flgContext) ; } template<typename TFunctional> requires(std::invocable<TFunctional, TDerived*> && std::tuple_size_v<AsTupleT<std::invoke_result_t<TFunctional, TDerived*>>> > 0) auto Chain(TFunctional&& fnTask, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { using ResultT = AsTupleT<std::invoke_result_t<TFunctional, TDerived*>>; using ResultVariantT = std::variant<ResultT, VoidT>; auto svrResult = ResourcePolicyT::template AllocateShared<ResultVariantT>(); return SequenceChainerT<true, ResultT> { this->Self(), svrResult, !(flgContext & Flags::Deferred), Enqueue( [fnTask = std::move(fnTask), svrResult] (TDerived* pSelf) mutable { *svrResult = std::invoke(std::move(fnTask), pSelf); }, idQueue, flgContext ) }; } template<typename TFunctional> requires(std::invocable<TFunctional, TDerived*> && std::tuple_size_v<AsTupleT<std::invoke_result_t<TFunctional, TDerived*>>> == 0) auto Chain(TFunctional&& fnTask, uint32_t idQueue = PolicyT::DefaultQueue(), Flags flgContext = Flags::Undefined) { return SequenceChainerT<true> { this->Self(), !(flgContext & Flags::Deferred), Enqueue(std::forward<TFunctional>(fnTask), idQueue, flgContext) }; } auto Chain(Flags flgContext = Flags::Undefined) { return SequenceChainerT<true> { this->Self(), !(flgContext & Flags::Deferred), true }; } ////////////////////////////////////////////////////////////////////////////////////////////////////////////////////////// size_t GetContentionLevel() const noexcept { if constexpr(PolicyT::TrackContention()) return nContention_.load(std::memory_order::relaxed); else return 0; } void ResetContentionLevel() noexcept { if constexpr(PolicyT::TrackContention()) nContention_.store(0, std::memory_order::relaxed); } protected: friend class Factory<TDerived>; using AtomicWaitableCounterT = std::atomic<CompatibleCounterT>; using AtomicCounterT = std::atomic<size_t>; static void DecCounter(AtomicCounterT* pCounter) noexcept { pCounter->fetch_sub(1, std::memory_order::relaxed); } using CounterGuardT = std::unique_ptr<AtomicCounterT, std::integral_constant<decltype(&DecCounter), &DecCounter>>; template<std::invocable<TDerived*> TFunctional> bool OnEnqueue(TFunctional&& /*fnTask*/, uint32_t /*idQueue*/, Flags /*flgContext*/) { return false; } template<std::invocable<TDerived*> TFunctional> std::pair<bool, uint32_t> OnComplete(Flags /*flgContext*/, TFunctional&& /*onComplete*/) { return { false, 0 }; } template<bool t_bNotify = true> std::pair<bool, size_t> OnEnqueued(uint32_t) noexcept { auto const nEnqueued = nEnqueued_.fetch_add(1, std::memory_order::release); if constexpr(t_bNotify) if(nEnqueued == 0) nEnqueued_.notify_one(); return { true, nEnqueued }; } template<bool t_bProcessed = false> size_t OnDequeued(uint32_t, CompatibleCounterT nDequeue = 1) noexcept { return nEnqueued_.fetch_sub(nDequeue, std::memory_order::relaxed); } size_t OnProcess(bool, uint32_t) { return 0; } bool OnSubmit() { return nEnqueued_.load(std::memory_order::acquire) != 0 && ExecutableT::OnSubmit() ; } void OnFinalize() noexcept { nEnqueued_.store(0, std::memory_order::release); ExecutableT::OnFinalize(); } bool HasEnqueued() const noexcept { return nEnqueued_.load(std::memory_order::acquire) != 0; } void WaitEnqueue(CompatibleCounterT nWhile = 0) const noexcept { nEnqueued_.wait(nWhile, std::memory_order::acquire); } FFPP_ATTR_INLINE CounterGuardT BeginTrackContention() { CounterGuardT sgAccessCounter; if constexpr(PolicyT::TrackContention()) { sgAccessCounter.reset(&nAccess_); auto const nAccess = nAccess_.fetch_add(1, std::memory_order::relaxed); auto nContention = nContention_.load(std::memory_order::relaxed); while( nAccess > nContention && !nContention_.compare_exchange_weak( nContention, nAccess, std::memory_order::relaxed, std::memory_order::relaxed ) ); } return sgAccessCounter; } protected: AtomicWaitableCounterT //Task counter. Immediate value represents current number of pending tasks. Can be used to wait for task (WaitEnqueue). nEnqueued_ alignas(PolicyT::ResourcePolicyT::InterferenceSize()) = 0; AtomicCounterT //Pair of counters for contention estimation (concurrent Enqueue method invocations, not used by default). nAccess_ alignas(PolicyT::ResourcePolicyT::InterferenceSize()) = 0 , nContention_ = 0 ; };//FunctionalSequence }//FFPP #endif//FFPP_FUNCTIONAL_SEQUENCE_HPP