#include "vdm/util/event_bus.hpp" #include #include #include #include #include "vtest.hpp" using vdm::EventBus; namespace { struct Progress { int task; long downloaded; }; struct StateChange { int task; std::string state; }; } // namespace VT_TEST(bus_delivers_to_matching_type_only) { EventBus bus; int progress_hits = 0; int state_hits = 0; bus.subscribe([&](const Progress &p) { ++progress_hits; VT_CHECK_EQ(p.task, 7); }); bus.subscribe([&](const StateChange &) { ++state_hits; }); bus.publish(Progress{7, 1024}); VT_CHECK_EQ(progress_hits, 1); VT_CHECK_EQ(state_hits, 0); bus.publish(StateChange{7, "downloading"}); VT_CHECK_EQ(progress_hits, 1); VT_CHECK_EQ(state_hits, 1); } VT_TEST(bus_invokes_in_registration_order) { EventBus bus; std::vector order; bus.subscribe([&](const Progress &) { order.push_back(1); }); bus.subscribe([&](const Progress &) { order.push_back(2); }); bus.subscribe([&](const Progress &) { order.push_back(3); }); bus.publish(Progress{0, 0}); VT_REQUIRE(order.size() == 3); VT_CHECK_EQ(order[0], 1); VT_CHECK_EQ(order[1], 2); VT_CHECK_EQ(order[2], 3); } VT_TEST(bus_unsubscribe_stops_delivery) { EventBus bus; int hits = 0; auto tok = bus.subscribe([&](const Progress &) { ++hits; }); bus.publish(Progress{0, 0}); bus.unsubscribe(tok); bus.publish(Progress{0, 0}); VT_CHECK_EQ(hits, 1); } VT_TEST(bus_scoped_subscription_auto_unsubscribes) { EventBus bus; int hits = 0; { auto sub = bus.subscribe_scoped([&](const Progress &) { ++hits; }); bus.publish(Progress{0, 0}); } bus.publish(Progress{0, 0}); VT_CHECK_EQ(hits, 1); } VT_TEST(bus_handler_may_unsubscribe_itself_during_dispatch) { EventBus bus; int hits = 0; EventBus::Token tok = EventBus::kInvalid; tok = bus.subscribe([&](const Progress &) { ++hits; bus.unsubscribe(tok); // must not deadlock or invalidate the dispatch loop }); bus.publish(Progress{0, 0}); bus.publish(Progress{0, 0}); VT_CHECK_EQ(hits, 1); } VT_TEST(bus_concurrent_publish_is_safe) { EventBus bus; std::atomic total{0}; bus.subscribe([&](const Progress &p) { total += p.downloaded; }); constexpr int kThreads = 8; constexpr int kPerThread = 2000; std::vector ts; for (int i = 0; i < kThreads; ++i) ts.emplace_back([&, i] { for (int j = 0; j < kPerThread; ++j) bus.publish(Progress{i, 1}); }); ts.clear(); // join VT_CHECK_EQ(total.load(), static_cast(kThreads) * kPerThread); }