#include #include #include #include #include Y_UNIT_TEST_SUITE(HugeCluster) { using namespace NActors; class TPoller: public TActor { const std::vector& Targets; std::unordered_map& Connected; public: TPoller(const std::vector& targets, std::unordered_map& events) : TActor(&TPoller::StateFunc) , Targets(targets) , Connected(events) {} void Handle(TEvTestStartPolling::TPtr /*ev*/, const TActorContext& ctx) { for (ui32 i = 0; i < Targets.size(); ++i) { ctx.Send(Targets[i], new TEvTest(), IEventHandle::FlagTrackDelivery, i); } } void Handle(TEvents::TEvUndelivered::TPtr ev, const TActorContext& ctx) { const ui32 cookie = ev->Cookie; // Cerr << "TEvUndelivered ping from node# " << SelfId().NodeId() << " to node# " << cookie + 1 << Endl; ctx.Send(Targets[cookie], new TEvTest(), IEventHandle::FlagTrackDelivery, cookie); } void Handle(TEvTest::TPtr ev, const TActorContext& /*ctx*/) { // Cerr << "Polled from " << ev->Sender.ToString() << Endl; Connected[ev->Sender].Signal(); } void Handle(TEvents::TEvPoisonPill::TPtr& /*ev*/, const TActorContext& ctx) { Die(ctx); } STRICT_STFUNC(StateFunc, HFunc(TEvents::TEvUndelivered, Handle) HFunc(TEvTestStartPolling, Handle) HFunc(TEvTest, Handle) HFunc(TEvents::TEvPoisonPill, Handle) ) }; class TStartPollers : public TActorBootstrapped { const std::vector& Pollers; public: TStartPollers(const std::vector& pollers) : Pollers(pollers) {} void Bootstrap(const TActorContext& ctx) { Become(&TThis::StateFunc); for (ui32 i = 0; i < Pollers.size(); ++i) { ctx.Send(Pollers[i], new TEvTestStartPolling(), IEventHandle::FlagTrackDelivery, i); } } void Handle(TEvents::TEvUndelivered::TPtr ev, const TActorContext& ctx) { const ui32 cookie = ev->Cookie; // Cerr << "TEvUndelivered start poller message to node# " << cookie + 1 << Endl; ctx.Send(Pollers[cookie], new TEvTestStartPolling(), IEventHandle::FlagTrackDelivery, cookie); } void Handle(TEvents::TEvPoisonPill::TPtr& /*ev*/, const TActorContext& ctx) { Die(ctx); } STRICT_STFUNC(StateFunc, HFunc(TEvents::TEvUndelivered, Handle) HFunc(TEvents::TEvPoisonPill, Handle) ) }; TIntrusivePtr MakeLogConfigs(NLog::EPriority priority) { // custom logger settings auto loggerSettings = MakeIntrusive( TActorId(0, "logger"), (NLog::EComponent)410, priority, priority, 0U); loggerSettings->Append( NActorsServices::EServiceCommon_MIN, NActorsServices::EServiceCommon_MAX, NActorsServices::EServiceCommon_Name ); constexpr ui32 WilsonComponentId = 430; // NKikimrServices::WILSON static const TString WilsonComponentName = "WILSON"; loggerSettings->Append( (NLog::EComponent)WilsonComponentId, (NLog::EComponent)WilsonComponentId + 1, [](NLog::EComponent) -> const TString & { return WilsonComponentName; }); return loggerSettings; } Y_UNIT_TEST(AllToAll) { ui32 nodesNum = 120; std::vector pollers(nodesNum); std::vector> events(nodesNum); // Must destroy actor system before shared arrays { TTestICCluster testCluster(nodesNum, NActors::TChannelsConfig(), nullptr, MakeLogConfigs(NLog::PRI_EMERG)); for (ui32 i = 0; i < nodesNum; ++i) { pollers[i] = testCluster.RegisterActor(new TPoller(pollers, events[i]), i + 1); } for (ui32 i = 0; i < nodesNum; ++i) { for (const auto& actor : pollers) { events[i][actor] = TManualEvent(); } } testCluster.RegisterActor(new TStartPollers(pollers), 1); for (ui32 i = 0; i < nodesNum; ++i) { for (auto& [_, ev] : events[i]) { ev.WaitI(); } } } } Y_UNIT_TEST(AllToOne) { ui32 nodesNum = 500; std::vector listeners; std::vector pollers(nodesNum - 1); std::unordered_map events; std::unordered_map emptyEventList; // Must destroy actor system before shared arrays { TTestICCluster testCluster(nodesNum, NActors::TChannelsConfig(), nullptr, MakeLogConfigs(NLog::PRI_EMERG)); const TActorId listener = testCluster.RegisterActor(new TPoller({}, events), nodesNum); listeners = { listener }; for (ui32 i = 0; i < nodesNum - 1; ++i) { pollers[i] = testCluster.RegisterActor(new TPoller(listeners, emptyEventList), i + 1); } for (const auto& actor : pollers) { events[actor] = TManualEvent(); } testCluster.RegisterActor(new TStartPollers(pollers), 1); for (auto& [_, ev] : events) { ev.WaitI(); } } } }