#include #include #include #include #include #include #include #include #include #include #include "renderive/base/observer/Observer.hpp" #include "renderive/frame_control/Frame_Control.hpp" #include "renderive/real_time_data/Real_Time_Data.hpp" #include "renderive/renderable/Renderable.hpp" #include "renderive/scene/Scene.hpp" #include "renderive/state/Double_State_Strategy.hpp" #include "renderive/state/Triple_State_Strategy.hpp" namespace { struct Threading_Test_Frame { std::uint64_t value{}; std::uint64_t mirror{~std::uint64_t{}}; }; bool valid_frame(const Threading_Test_Frame& frame) { return frame.mirror == ~frame.value; } void wait_start(const std::atomic& start) { while (!start.load(std::memory_order_acquire)) { std::this_thread::yield(); } } struct Threading_Test_Observation {}; struct Threading_Test_Observer_Data { std::atomic active{}; std::atomic maximum{}; std::atomic count{}; }; struct Threading_Test_Observer { static constexpr bool enabled = true; std::shared_ptr data{std::make_shared()}; void observe(const Threading_Test_Observation&) noexcept { const int active = data->active.fetch_add(1, std::memory_order_acq_rel) + 1; int maximum = data->maximum.load(std::memory_order_acquire); while (maximum < active && !data->maximum.compare_exchange_weak(maximum, active, std::memory_order_acq_rel)) {} std::this_thread::yield(); data->count.fetch_add(1, std::memory_order_relaxed); data->active.fetch_sub(1, std::memory_order_acq_rel); } }; struct Threading_Test_Time_Source { std::uint64_t now_ns() const noexcept { return 0; } }; using Threading_Test_Observer_State = Observer_State; struct Threading_Revision_Observer_Data { std::mutex mutex; std::vector revisions; }; struct Threading_Revision_Observer { static constexpr bool enabled = true; std::shared_ptr data{std::make_shared()}; void observe(const Real_Time_Data_Observation& observation) noexcept { std::lock_guard lock(data->mutex); data->revisions.push_back(observation.state.revision); } }; using Threading_Revision_Observer_State = Observer_State; struct Threading_Test_State_Base {}; struct Threading_Test_State_Value { Threading_Test_Frame frame; }; using Threading_Test_Double_State = Double_State_Strategy; using Threading_Test_Triple_State = Triple_State_Strategy; struct Threading_Test_Scene_Renderable : Renderable_Base { Threading_Test_Scene_Renderable(std::atomic& render_count, std::atomic& maximum_sequence) : Renderable_Base({.cache_enabled = false}), render_count(&render_count), maximum_sequence(&maximum_sequence) {} void prepare(const Prepare_Render_Context& context) override { render_count->fetch_add(1, std::memory_order_relaxed); std::uint64_t maximum = maximum_sequence->load(std::memory_order_acquire); while (maximum < context.frame.render_sequence && !maximum_sequence->compare_exchange_weak(maximum, context.frame.render_sequence, std::memory_order_acq_rel)) {} } std::atomic* render_count; std::atomic* maximum_sequence; }; struct Threading_Test_Render_Graph_Renderable : Renderable_Base { Threading_Test_Render_Graph_Renderable(std::atomic& executed) : Renderable_Base({.cache_enabled = false}), executed(&executed) {} void build_prepare_graph(Renderable_Graph_Builder& graph) override { graph.emplace("prepare", "Prepare", [this](const Prepare_Render_Context&) { executed->fetch_add(1, std::memory_order_relaxed); }); } void request_graph_rebuild() { rebuild_render_graph(); } std::atomic* executed; }; } TEST(threading_contract_test, observer_state_serializes_callbacks_from_multiple_threads) { Threading_Test_Observer observer; auto data = observer.data; Threading_Test_Observer_State state(observer, Threading_Test_Time_Source{}); constexpr int thread_count = 8; constexpr int observations_per_thread = 200; std::atomic start{}; std::vector threads; threads.reserve(thread_count); for (int index = 0; index < thread_count; ++index) { threads.emplace_back([&] { wait_start(start); for (int observation = 0; observation < observations_per_thread; ++observation) { state.observe(Threading_Test_Observation{}); } }); } start.store(true, std::memory_order_release); for (auto& thread : threads) { thread.join(); } EXPECT_EQ(data->maximum.load(std::memory_order_acquire), 1); EXPECT_EQ(data->count.load(std::memory_order_acquire), thread_count * observations_per_thread); } TEST(threading_contract_test, flow_supports_multiple_producers_and_multiple_consumers) { Flow_Refresh_Strategy strategy; constexpr int producer_count = 4; constexpr int consumer_count = 4; constexpr int frames_per_producer = 500; constexpr int frame_count = producer_count * frames_per_producer; auto seen = std::make_unique[]>(frame_count); std::atomic start{}; std::atomic producers_done{}; std::atomic consumed{}; std::atomic invalid{}; std::vector consumers; consumers.reserve(consumer_count); for (int index = 0; index < consumer_count; ++index) { consumers.emplace_back([&] { wait_start(start); for (;;) { auto frame = strategy.acquire_renderer(); if (frame) { const auto value = frame->value; if (!valid_frame(*frame) || value >= static_cast(frame_count)) { invalid.fetch_add(1, std::memory_order_relaxed); } else if (seen[value].fetch_add(1, std::memory_order_acq_rel) != 0) { invalid.fetch_add(1, std::memory_order_relaxed); } consumed.fetch_add(1, std::memory_order_release); continue; } if (producers_done.load(std::memory_order_acquire) == producer_count && strategy.pending_frame_count() == 0) { break; } std::this_thread::yield(); } }); } std::vector producers; producers.reserve(producer_count); for (int producer = 0; producer < producer_count; ++producer) { producers.emplace_back([&, producer] { wait_start(start); for (int index = 0; index < frames_per_producer; ++index) { const auto value = static_cast(producer * frames_per_producer + index); auto frame = strategy.acquire_painter(); frame->value = value; frame->mirror = ~value; } producers_done.fetch_add(1, std::memory_order_release); }); } start.store(true, std::memory_order_release); for (auto& producer : producers) { producer.join(); } for (auto& consumer : consumers) { consumer.join(); } EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(consumed.load(std::memory_order_acquire), frame_count); EXPECT_EQ(strategy.pending_frame_count(), 0); const auto state = strategy.state(); EXPECT_EQ(state.enqueued_frame_count, frame_count); EXPECT_EQ(state.dequeued_frame_count, frame_count); EXPECT_EQ(state.rendered_frame_count, frame_count); } TEST(threading_contract_test, manual_strategy_keeps_frames_consistent_during_concurrent_paint_refresh_and_render) { Manual_Refresh_Strategy strategy; constexpr int producer_count = 4; constexpr int frames_per_producer = 500; constexpr int frame_count = producer_count * frames_per_producer; std::atomic start{}; std::atomic producers_done{}; std::atomic refresher_done{}; std::atomic invalid{}; std::vector producers; producers.reserve(producer_count); for (int producer = 0; producer < producer_count; ++producer) { producers.emplace_back([&, producer] { wait_start(start); for (int index = 0; index < frames_per_producer; ++index) { const auto value = static_cast(producer * frames_per_producer + index + 1); auto frame = strategy.acquire_painter(); frame->value = value; frame->mirror = ~value; } producers_done.fetch_add(1, std::memory_order_release); }); } std::thread refresher([&] { wait_start(start); for (;;) { strategy.refresh(); const auto state = strategy.state(); if (producers_done.load(std::memory_order_acquire) == producer_count && !state.pending_frame) { break; } std::this_thread::yield(); } refresher_done.store(true, std::memory_order_release); }); std::thread renderer([&] { wait_start(start); while (!refresher_done.load(std::memory_order_acquire)) { auto frame = strategy.acquire_renderer(); if (frame && !valid_frame(*frame)) { invalid.fetch_add(1, std::memory_order_relaxed); } std::this_thread::yield(); } auto frame = strategy.acquire_renderer(); if (frame && !valid_frame(*frame)) { invalid.fetch_add(1, std::memory_order_relaxed); } }); start.store(true, std::memory_order_release); for (auto& producer : producers) { producer.join(); } refresher.join(); renderer.join(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); const auto state = strategy.state(); EXPECT_EQ(state.prepared_frame_count, frame_count); EXPECT_GT(state.successful_refresh_count, 0); EXPECT_FALSE(state.pending_frame); } TEST(threading_contract_test, low_latency_strategy_keeps_frames_consistent_during_concurrent_operations) { Low_Latency_Strategy strategy; constexpr int producer_count = 4; constexpr int frames_per_producer = 500; constexpr int frame_count = producer_count * frames_per_producer; constexpr int update_count = 2000; std::atomic start{}; std::atomic producers_done{}; std::atomic updater_done{}; std::atomic configuration_done{}; std::atomic invalid{}; std::vector producers; producers.reserve(producer_count); for (int producer = 0; producer < producer_count; ++producer) { producers.emplace_back([&, producer] { wait_start(start); for (int index = 0; index < frames_per_producer; ++index) { const auto value = static_cast(producer * frames_per_producer + index + 1); auto frame = strategy.acquire_painter(); frame->value = value; frame->mirror = ~value; } producers_done.fetch_add(1, std::memory_order_release); }); } std::thread updater([&] { wait_start(start); for (int index = 1; index <= update_count; ++index) { strategy.on_real_time_data_update({Real_Time_Data_Observation_Event::updated, {nullptr, Real_Time_Data_Retention::latest, static_cast(index), static_cast(index), static_cast(index), 1}}); } updater_done.store(true, std::memory_order_release); }); std::thread configurator([&] { wait_start(start); constexpr double frequencies[] = {30.0, 60.0, 120.0, 240.0}; for (int index = 0; index < frame_count; ++index) { strategy.set_frequency_hz(frequencies[index % 4]); } configuration_done.store(true, std::memory_order_release); }); std::thread renderer([&] { wait_start(start); while (producers_done.load(std::memory_order_acquire) != producer_count || !updater_done.load(std::memory_order_acquire) || !configuration_done.load(std::memory_order_acquire)) { auto frame = strategy.acquire_renderer(); if (frame && !valid_frame(*frame)) { invalid.fetch_add(1, std::memory_order_relaxed); } std::this_thread::yield(); } for (int index = 0; index < 100; ++index) { auto frame = strategy.acquire_renderer(); if (frame && !valid_frame(*frame)) { invalid.fetch_add(1, std::memory_order_relaxed); } } }); start.store(true, std::memory_order_release); for (auto& producer : producers) { producer.join(); } updater.join(); configurator.join(); renderer.join(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(strategy.state().real_time_data_update_sequence, update_count); EXPECT_TRUE(std::isfinite(strategy.state().frequency_hz)); } TEST(threading_contract_test, real_time_data_observations_follow_mutation_revision_order) { Threading_Revision_Observer observer; auto observer_data = observer.data; Latest_Real_Time_Data data(Threading_Revision_Observer_State(observer, Threading_Test_Time_Source{})); constexpr int updater_count = 8; constexpr int updates_per_thread = 500; constexpr int update_count = updater_count * updates_per_thread; std::atomic start{}; std::vector updaters; updaters.reserve(updater_count); for (int updater = 0; updater < updater_count; ++updater) { updaters.emplace_back([&, updater] { wait_start(start); for (int index = 0; index < updates_per_thread; ++index) { data.update(updater * updates_per_thread + index); } }); } start.store(true, std::memory_order_release); for (auto& updater : updaters) { updater.join(); } std::lock_guard lock(observer_data->mutex); EXPECT_EQ(observer_data->revisions.size(), update_count); if (observer_data->revisions.size() == update_count) { for (std::size_t index = 0; index < observer_data->revisions.size(); ++index) { EXPECT_EQ(observer_data->revisions[index], index + 1); } } } TEST(threading_contract_test, latest_real_time_data_supports_concurrent_updates_and_snapshots) { Latest_Real_Time_Data data; constexpr int updater_count = 4; constexpr int updates_per_thread = 1000; constexpr int update_count = updater_count * updates_per_thread; std::atomic start{}; std::atomic updates_done{}; std::atomic invalid{}; std::thread reader([&] { wait_start(start); while (updates_done.load(std::memory_order_acquire) != updater_count) { const auto snapshot = data.snapshot(); if (snapshot && !valid_frame(*snapshot)) { invalid.fetch_add(1, std::memory_order_relaxed); } const auto state = data.update_state(); if (state.retained_value_count > 1) { invalid.fetch_add(1, std::memory_order_relaxed); } } }); std::vector updaters; updaters.reserve(updater_count); for (int updater = 0; updater < updater_count; ++updater) { updaters.emplace_back([&, updater] { wait_start(start); for (int index = 0; index < updates_per_thread; ++index) { const auto value = static_cast(updater * updates_per_thread + index); data.update({value, ~value}); } updates_done.fetch_add(1, std::memory_order_release); }); } start.store(true, std::memory_order_release); for (auto& updater : updaters) { updater.join(); } reader.join(); const auto state = data.update_state(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(state.revision, update_count); EXPECT_EQ(state.total_update_count, update_count); EXPECT_EQ(state.retained_value_count, 1); } TEST(threading_contract_test, history_real_time_data_supports_concurrent_update_snapshot_and_discard) { History_Real_Time_Data data; constexpr int updater_count = 4; constexpr int updates_per_thread = 500; constexpr int update_count = updater_count * updates_per_thread; std::atomic start{}; std::atomic updates_done{}; std::atomic discard_done{}; std::atomic invalid{}; std::thread reader([&] { wait_start(start); while (!discard_done.load(std::memory_order_acquire)) { const auto snapshot = data.snapshot(); for (const auto& frame : snapshot) { if (!valid_frame(frame)) { invalid.fetch_add(1, std::memory_order_relaxed); break; } } } }); std::thread discarder([&] { wait_start(start); while (updates_done.load(std::memory_order_acquire) != updater_count) { data.discard_before_time_ns(std::numeric_limits::max()); std::this_thread::yield(); } data.discard_before_time_ns(std::numeric_limits::max()); discard_done.store(true, std::memory_order_release); }); std::vector updaters; updaters.reserve(updater_count); for (int updater = 0; updater < updater_count; ++updater) { updaters.emplace_back([&, updater] { wait_start(start); for (int index = 0; index < updates_per_thread; ++index) { const auto value = static_cast(updater * updates_per_thread + index); data.update({value, ~value}); } updates_done.fetch_add(1, std::memory_order_release); }); } start.store(true, std::memory_order_release); for (auto& updater : updaters) { updater.join(); } discarder.join(); reader.join(); const auto state = data.update_state(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(state.total_update_count, update_count); EXPECT_EQ(state.retained_value_count, data.size()); EXPECT_GE(state.revision, update_count); } TEST(threading_contract_test, double_state_keeps_cache_consistent_during_concurrent_set_publish_and_read) { Threading_Test_Double_State strategy; constexpr int update_count = 4000; std::atomic start{}; std::atomic writer_done{}; std::atomic invalid{}; std::thread writer([&] { wait_start(start); for (int index = 1; index <= update_count; ++index) { const auto value = static_cast(index); strategy.set<&Threading_Test_State_Value::frame>(Threading_Test_Frame{value, ~value}); strategy.publish(); } writer_done.store(true, std::memory_order_release); }); std::thread cache_reader([&] { wait_start(start); while (!writer_done.load(std::memory_order_acquire)) { if (!valid_frame(strategy.get<&Threading_Test_State_Value::frame>())) { invalid.fetch_add(1, std::memory_order_relaxed); } } }); start.store(true, std::memory_order_release); writer.join(); cache_reader.join(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(strategy.state_revision(), update_count); } TEST(threading_contract_test, triple_state_returns_consistent_snapshots_during_concurrent_publish_and_acquire) { Threading_Test_Triple_State strategy; constexpr int update_count = 4000; std::atomic start{}; std::atomic writer_done{}; std::atomic invalid{}; std::thread writer([&] { wait_start(start); for (int index = 1; index <= update_count; ++index) { const auto value = static_cast(index); strategy.set<&Threading_Test_State_Value::frame>(Threading_Test_Frame{value, ~value}); strategy.publish(); } writer_done.store(true, std::memory_order_release); }); std::thread renderer([&] { wait_start(start); while (!writer_done.load(std::memory_order_acquire)) { strategy.acquire_render_state(); } strategy.acquire_render_state(); }); std::thread cache_reader([&] { wait_start(start); while (!writer_done.load(std::memory_order_acquire)) { if (!valid_frame(strategy.get<&Threading_Test_State_Value::frame>())) { invalid.fetch_add(1, std::memory_order_relaxed); } } }); start.store(true, std::memory_order_release); writer.join(); renderer.join(); cache_reader.join(); EXPECT_EQ(invalid.load(std::memory_order_acquire), 0); EXPECT_EQ(strategy.state_revision(), update_count); } TEST(threading_contract_test, scene_serializes_multiple_render_submitters_without_losing_submissions) { Scene2D_Context<> scene; constexpr int submitter_count = 4; constexpr int renders_per_submitter = 100; constexpr int render_count = submitter_count * renders_per_submitter; std::atomic executed{}; std::atomic maximum_sequence{}; auto renderable = renderive_Owner::make(executed, maximum_sequence); { auto attach = scene.attach_builder(); attach.attach(renderable); } std::atomic start{}; std::vector submitters; submitters.reserve(submitter_count); for (int index = 0; index < submitter_count; ++index) { submitters.emplace_back([&] { wait_start(start); for (int render = 0; render < renders_per_submitter; ++render) { scene.render(); scene.wait_for_render(); } }); } start.store(true, std::memory_order_release); for (auto& submitter : submitters) { submitter.join(); } EXPECT_EQ(executed.load(std::memory_order_acquire), render_count); EXPECT_EQ(maximum_sequence.load(std::memory_order_acquire), render_count); } TEST(threading_contract_test, scene_render_and_render_graph_rebuild_are_serialized_by_mutation_barrier) { Scene2D_Context<> scene; constexpr int render_count = 300; std::atomic executed{}; auto renderable = renderive_Owner::make(executed); { auto attach = scene.attach_builder(); attach.attach(renderable); } std::atomic start{}; std::thread renderer([&] { wait_start(start); for (int index = 0; index < render_count; ++index) { scene.render(); scene.wait_for_render(); } }); std::thread rebuilder([&] { wait_start(start); for (int index = 0; index < render_count; ++index) renderable->request_graph_rebuild(); }); start.store(true, std::memory_order_release); renderer.join(); rebuilder.join(); scene.render(); scene.wait_for_render(); EXPECT_EQ(executed.load(std::memory_order_acquire), render_count + 1); } TEST(threading_contract_test, real_time_data_updates_can_race_with_scene_destruction) { using Data = Latest_Real_Time_Data>; struct Renderable : Renderable_Base { explicit Renderable() : Renderable_Base() {} }; using Attachment = Attach_Real_Time_Data; auto data = std::make_shared(); auto scene = std::make_unique>(); auto renderable = renderive_Owner::make(With_Real_Time_Data(data)); { auto attach = scene->attach_builder(); attach.attach(renderable); } std::atomic start{}; std::atomic stop{}; std::atomic updates{}; std::thread updater([&] { wait_start(start); int value{}; while (!stop.load(std::memory_order_acquire)) { data->update(++value); updates.fetch_add(1, std::memory_order_relaxed); } }); start.store(true, std::memory_order_release); while (updates.load(std::memory_order_acquire) < 100) { std::this_thread::yield(); } scene.reset(); while (updates.load(std::memory_order_acquire) < 200) { std::this_thread::yield(); } stop.store(true, std::memory_order_release); updater.join(); EXPECT_GE(data->revision(), 200); renderable.reset(); } TEST(threading_contract_test, frame_leases_explicitly_declare_same_thread_lifetime) { using Flow = Flow_Refresh_Strategy; using Manual = Manual_Refresh_Strategy; using Low_Latency = Low_Latency_Strategy; EXPECT_TRUE(Flow::Painter_Lease::thread_affine); EXPECT_TRUE(Flow::Render_Lease::thread_affine); EXPECT_TRUE(Manual::Painter_Lease::thread_affine); EXPECT_TRUE(Manual::Render_Lease::thread_affine); EXPECT_TRUE(Low_Latency::Painter_Lease::thread_affine); EXPECT_TRUE(Low_Latency::Render_Lease::thread_affine); }