/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
include/actor.hpp
245 строк
8 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_ACTOR_HPP #define FFPP_ACTOR_HPP namespace FFPP::Concepts { template<typename TPool, typename TActor> concept ActorPool = Executor<TPool, TActor>; }//FFPP::Concepts namespace FFPP::Base { template<typename TDerived, Concepts::FunctionalQueuePolicy TPolicy> class Actor : public FunctionalQueue<TDerived, TPolicy> { public: using PolicyT = TPolicy; using ResourcePolicyT = PolicyT::ResourcePolicyT; template<typename Type> using SharedT = ResourcePolicyT::template SharedT<Type>; using FunctionalQueueT = FunctionalQueue<TDerived, TPolicy>; friend FunctionalQueueT; using FunctionalSequenceT = FunctionalQueueT::FunctionalSequenceT; friend FunctionalSequenceT; using ExecutableT = FunctionalSequenceT::ExecutableT; friend ExecutableT; template<typename TSharedExecutor> requires(Concepts::SharedExecutor<TSharedExecutor, TDerived>) Actor( TSharedExecutor const& spExecutor , uint32_t nPriority = PolicyT::DefaultPriority() , size_t nLimit = FunctionalSequenceT::InvalidLimit() #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) : FunctionalQueueT( nPriority , nLimit #if FFPP_TRACK_ORIGIN , slOrigin #endif ) , c_fnSubmitToExecutor_([ this , spExecutor = ValidateExecutor( spExecutor #if FFPP_TRACK_ORIGIN , slOrigin #endif ) ] () { return spExecutor->Submit(this->Self()); }) { } Actor(Actor const&) = delete; Actor(Actor&&) = delete; Actor& operator = (Actor const&) = delete; Actor& operator = (Actor&&) = delete; template<typename TSharedExecutor> requires(Concepts::SharedExecutor<TSharedExecutor, TDerived>) FFPP_ATTR_INLINE static SharedT<TDerived> MakeShared( TSharedExecutor const& spExecutor , uint32_t nPriority = PolicyT::DefaultPriority() , size_t nLimit = FunctionalSequenceT::InvalidLimit() #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin = std::source_location::current() #endif ) { return ResourcePolicyT::template AllocateShared<TDerived>( spExecutor , nPriority , nLimit #if FFPP_TRACK_ORIGIN , slOrigin #endif ); } protected: template<typename TSharedExecutor> static TSharedExecutor ValidateExecutor( TSharedExecutor const& spExecutor #if FFPP_TRACK_ORIGIN , std::source_location const& slOrigin #endif ) { if(spExecutor == nullptr) Except<std::invalid_argument>::Throw( "Null actor executor pointer @" #if FFPP_TRACK_ORIGIN , slOrigin #endif ); return spExecutor; } template<bool t_bProcessed = false> size_t OnDequeued(uint32_t idQueue, CompatibleCounterT nDequeue = 1) noexcept { auto const nEnqueued = FunctionalQueueT::template OnDequeued<t_bProcessed>(idQueue, nDequeue); if(nEnqueued == 1) this->UpdateProcessingThread({ }); return nEnqueued; } bool SubmitToExecutor() const { return std::invoke(c_fnSubmitToExecutor_); } private: ResourcePolicyT::template FunctionT<bool(), false> c_fnSubmitToExecutor_; };//Actor }//FFPP::Base namespace FFPP { template<typename TDerived, bool t_bImmediate = true, Concepts::FunctionalQueuePolicy TPolicy = FunctionalQueuePolicy> class Actor : public Base::Actor<TDerived, TPolicy> { public: using ActorBaseT = Base::Actor<TDerived, TPolicy>; using PolicyT = ActorBaseT::PolicyT; using ResourcePolicyT = ActorBaseT::ResourcePolicyT; using FunctionalQueueT = ActorBaseT::FunctionalQueueT; friend FunctionalQueueT; using FunctionalSequenceT = FunctionalQueueT::FunctionalSequenceT; friend FunctionalSequenceT; using ExecutableT = FunctionalSequenceT::ExecutableT; friend ExecutableT; 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>; using ActorBaseT::ActorBaseT; size_t Process([[maybe_unused]] bool bAll = false, [[maybe_unused]] uint32_t idQueue = FunctionalQueueT::InvalidQueueId()) { this->UpdateProcessingThread(); return FunctionalQueueT::OnProcessSingle(); } protected: template<bool t_bNotify = false> //Actors's enqueue counter is not used for any thread notification. std::pair<bool, size_t> OnEnqueued(uint32_t idQueue) { auto [bEnqueued, nEnqueued] = ActorBaseT::template OnEnqueued<t_bNotify>(idQueue); if(!bEnqueued) return { false, nEnqueued }; if constexpr(TestModeFlags(PolicyT::SchedulingMethod(), SchedulingFlags::Deferred)) { if(nEnqueued == 0) bEnqueued = this->SubmitToExecutor(); } else { bEnqueued = this->SubmitToExecutor(); } return { bEnqueued, nEnqueued }; } template<bool t_bProcessed = false> size_t OnDequeued(uint32_t idQueue, CompatibleCounterT nDequeue = 1) noexcept(!t_bProcessed) { auto const nEnqueued = ActorBaseT::template OnDequeued<t_bProcessed>(idQueue, nDequeue); //Actor must not be submitted to thread pool during cleanup in deferred scheduling mode. if constexpr(TestModeFlags(PolicyT::SchedulingMethod(), SchedulingFlags::Deferred) && t_bProcessed) { if(nEnqueued > 1) { if(!this->SubmitToExecutor()) { Except<>::Throw( "Failed to submit Actor instance to executor @" #if FFPP_TRACK_ORIGIN , this->GetOrigin() #endif ); } } } return nEnqueued; } };//Actor template<typename TDerived, Concepts::FunctionalQueuePolicy TPolicy> class Actor<TDerived, false, TPolicy> : public Base::Actor<TDerived, TPolicy> { public: using ActorBaseT = Base::Actor<TDerived, TPolicy>; using ResourcePolicyT = ActorBaseT::ResourcePolicyT; using FunctionalQueueT = ActorBaseT::FunctionalQueueT; friend FunctionalQueueT; using FunctionalSequenceT = FunctionalQueueT::FunctionalSequenceT; friend FunctionalSequenceT; using ExecutableT = FunctionalSequenceT::ExecutableT; friend ExecutableT; 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>; using ActorBaseT::ActorBaseT; size_t Process([[maybe_unused]] bool bAll = true, [[maybe_unused]] uint32_t idQueue = FunctionalQueueT::InvalidQueueId()) { this->UpdateProcessingThread(); return FunctionalQueueT::OnProcessAll(); } protected: bool OnSubmit() { if(!FunctionalQueueT::OnSubmit()) return false; return this->SubmitToExecutor(); } };//Actor template<bool t_bImmediate = true, Concepts::FunctionalQueuePolicy TPolicy = FunctionalQueuePolicy> class FunctionalActor : public Actor<FunctionalActor<t_bImmediate, TPolicy>, t_bImmediate, TPolicy> { public: using ActorT = Actor<FunctionalActor<t_bImmediate, TPolicy>, t_bImmediate, TPolicy>; using FunctionalActorT = FunctionalActor<t_bImmediate, TPolicy>; using ResourcePolicyT = ActorT::ResourcePolicyT; using FunctionalQueueT = ActorT::FunctionalQueueT; using FunctionalSequenceT = FunctionalQueueT::FunctionalSequenceT; using ExecutableT = FunctionalSequenceT::ExecutableT; friend ExecutableT; template<typename Type> using SharedT = ResourcePolicyT::template SharedT<Type>; template<typename Type> using WeakT = ResourcePolicyT::template WeakT<Type>; using SP = SharedT<FunctionalActorT>; using WP = WeakT<FunctionalActorT>; using ActorT::ActorT; };//FunctionalActor }//FFPP #endif//FFPP_ACTOR_HPP