/
Ant010ff
/
ffpp
Обзор
Документация
Войти
/
Ant010ff
/
ffpp
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
test/actorring/main.cpp
126 строк
4 KB
Ant010ff
Demo code update; exception-related refactoring (FFPP_TRY / FFPP_CATCH).
13 июн 2026, 15:39
13 июн 2026, 15:39
d23feb5
Код
Авторство
О чём код?
/* Adoption of a Thread-ring test (https://salsa.debian.org/benchmarksgame-team/archive-alioth-benchmarksgame/-/tree/master/contributed-source-code/shootout/threadring) */ #include "../common/common.hpp" using ClockT = std::chrono::steady_clock; constexpr size_t c_nWorkers = 503; //Worker actors ("threads") in a ring. struct BasePolicy : CommonPolicy { using LockableT = ffpp::Lockable<ffpp::StripedTicketMutex<>>; static constexpr auto Queues() { return std::vector { 0 }; } static constexpr ffpp::DispatchFlags DispatchMethod() { //Special case - single-threaded pool - better suites this task. return ffpp::DispatchFlags::Single; } }; struct DequeActorPolicy : BasePolicy { using QueuePolicyT = ffpp::DequePolicy; template<typename TValue, typename TPolicy> using QueueT = ffpp::Deque<TValue, TPolicy>; }; struct RingQueueActorPolicy : BasePolicy { struct QueuePolicy : ffpp::RingQueuePolicy { using LockableT = BasePolicy::LockableT; static size_t QueueMinimalCapacity() { return c_nWorkers; } //queue minimal capacity, will be rounded to power of 2 }; using QueuePolicyT = QueuePolicy; }; struct FixedRingQueueActorPolicy : BasePolicy { struct QueuePolicy : ffpp::RingQueuePolicy { using LockableT = BasePolicy::LockableT; static constexpr auto QueueMode() { return ffpp::QueueFlags::Fixed; //fixed-capacity ring-queue } static size_t QueueCapacity() { return c_nWorkers; } //fixed queue capacity (if enabled) }; using QueuePolicyT = QueuePolicy; }; template<typename TPolicy> struct Worker : public ffpp::Actor<Worker<TPolicy>, true, TPolicy> { using WorkerT = Worker<TPolicy>; using WorkerVectorT = std::vector<std::shared_ptr<WorkerT>>; using ActorPoolT = ffpp::FunctionalQueuePool<TPolicy>; ClockT::time_point tpStop = ClockT::time_point::max(); static auto& GetActorPool() { static std::shared_ptr<ActorPoolT> s_spActorPool = std::make_shared<ActorPoolT>(); return s_spActorPool; } Worker() : ffpp::Actor<Worker<TPolicy>, true, TPolicy>(GetActorPool()) { } struct MsgWalkRing{ WorkerVectorT& vWorkers_; int iWorker_, nToken_; }; void On(MsgWalkRing& msg) { if(++msg.iWorker_ == msg.vWorkers_.size()) msg.iWorker_ = 0; if(--msg.nToken_) { msg.vWorkers_[msg.iWorker_]->Emit(msg); } else{ tpStop = ClockT::now(); Info::Put("Stoped at worker:", msg.iWorker_); for(auto const& spWorker : msg.vWorkers_) spWorker->Finalize(); GetActorPool()->Complete(); } } };//Worker template<typename TPolicy> auto RunTest(int nInitialToken) { using WorkerT = Worker<TPolicy>; typename WorkerT::WorkerVectorT vWorkers; vWorkers.reserve(c_nWorkers); for(size_t i = 0; i < c_nWorkers; ++i) vWorkers.push_back(std::make_shared<WorkerT>()); auto const tpStart = ClockT::now(); vWorkers[0]->Emit(typename WorkerT::MsgWalkRing { vWorkers, 0, nInitialToken }); WorkerT::GetActorPool()->Wait(); auto const tpStop = (*std::ranges::min_element(vWorkers, [] (auto const& left, auto const& right) { return left->tpStop < right->tpStop; }))->tpStop; auto const drnTest = tpStop - tpStart; Info::Put("Total process time:", ">>", std::chrono::duration_cast<std::chrono::milliseconds>(drnTest), "<<"); return drnTest; } int main([[maybe_unused]] int argc, [[maybe_unused]] char** argv) { FFPP_TRY { int const c_nInitialToken = argc > 1 ? std::stoi(argv[1]) : 5000000; Log<>::Init( algy::FilesystemOptions { algy::verbose, algy::c_bitDefaultCaps, std::filesystem::path(argv[0]).parent_path() / "log" , "actorring" }, algy::ConsoleOptions { algy::info, algy::c_bitDefaultCaps } ); Accent::Put("Actor-ring test started.", "Ring size:", c_nWorkers, ", number of iterations:", c_nInitialToken); Info::Put("Running exclusive deque test..."); RunTest<DequeActorPolicy>(c_nInitialToken); Info::Put("Running growing (default) ring queue test..."); auto drnGrowing = RunTest<RingQueueActorPolicy>(c_nInitialToken); Info::Put("Running fixed capacity ring queue test..."); auto drnFixed = RunTest<FixedRingQueueActorPolicy>(c_nInitialToken); Accent::Put("Growing vs fixed ring queue overhead:", int(100 * (double(drnGrowing.count()) / drnFixed.count() - 1)), "%"); return 0; } FFPP_CATCH(ffpp::Base::Exception, ex) { if(Log<>::IsValid()) Fatal::Put("Top level FFPP exception:", ex); return -2; } FFPP_CATCH(std::exception, ex) { if(Log<>::IsValid()) Fatal::Put("Top level exception:", ex); return -1; } }//main