diff --git a/Kernel/#U8bbe#U8ba1#U6587#U6863.md b/Kernel/#U8bbe#U8ba1#U6587#U6863.md new file mode 100644 index 0000000..5cf1a8e --- /dev/null +++ b/Kernel/#U8bbe#U8ba1#U6587#U6863.md @@ -0,0 +1,11 @@ +renderable 可以被渲染的基础元素 + +状态策略给类的状态提供线程安全访问更新机制 + +帧策略提供帧的发布生成机制 + +plot 上下文对象聚合所有策略 + +Flow 当前使用 Kernel 内部 `std::pmr::deque + mutex` 实现并发帧队列,不新增 Boost.Lockfree 依赖。多个 producer 和多个 consumer 调用的线程安全由测试覆盖,实际 renderer 消费通过 render mutex 串行化;Frame 对传入 PMR 的 allocate/deallocate 也在内部状态锁下串行化,所以允许使用 `std::pmr::unsynchronized_pool_resource`。 + +完整线程模型和多线程测试约束见 `threading.md`。 diff --git a/Kernel/CMakeLists.txt b/Kernel/CMakeLists.txt index a38a5ac..2d9ae9d 100644 --- a/Kernel/CMakeLists.txt +++ b/Kernel/CMakeLists.txt @@ -68,6 +68,7 @@ if (BUILD_TESTING) "${CMAKE_CURRENT_LIST_DIR}/tests/renderive/state/Concepts_Test.cpp" "${CMAKE_CURRENT_LIST_DIR}/tests/renderive/state/Double_State_Strategy_Test.cpp" "${CMAKE_CURRENT_LIST_DIR}/tests/renderive/state/Triple_State_Strategy_Test.cpp" + "${CMAKE_CURRENT_LIST_DIR}/tests/renderive/threading/Threading_Contract_Test.cpp" ) find_package(GTest CONFIG QUIET) if (NOT GTest_FOUND) diff --git a/Kernel/readme.md b/Kernel/readme.md index 46b59b0..9b9a82d 100644 --- a/Kernel/readme.md +++ b/Kernel/readme.md @@ -88,16 +88,20 @@ Taskflow 只在 `Scene_Base.cpp` 中包含。`Scene_Base` 通过内部执行上 `BUILD_TESTING=ON` 时构建 `renderive_scene_tests` 并通过 CTest 运行;`BUILD_TESTING=OFF` 时只构建 `renderive_scene` 库,不再编译 `tests/main.cpp` 或任何测试源文件。 +## 线程模型 + +Scene、Frame Strategy、Real-Time-Data、State Strategy、Renderable task graph 的允许并发组合、生命周期边界和对应测试统一记录在 `threading.md`。多线程语义测试作为普通测试套件的一部分构建,不依赖 ThreadSanitizer;Linux CI 建议额外用 TSan 执行同一套测试。 + ## Taskflow 查找 CMake 通过 `RENDERIVE_TASKFLOW_ROOT` 或系统 include 路径查找 `taskflow/taskflow.hpp`。找到 Taskflow 4.1.0 时使用真实执行器;找不到时使用 `src/renderive/compat/taskflow/taskflow.hpp` 的内置顺序依赖图执行器,使工程仍可完成 CMake 配置、编译并运行依赖图测试。兼容执行器只实现 Renderive 当前使用的 Taskflow 子集,用于验证内部 DAG、Renderable 外部依赖和 module 完成关系是否被正确转换,不模拟并行调度能力。真实 Taskflow 下额外运行并行执行测试。 -## Flow MPSC 队列 +## Flow 并发队列 -`Flow_Refresh_Strategy` 使用 `boost::lockfree::queue` 作为多生产者、单消费者帧队列。多个 painter 线程可以并发发布帧,renderer 按队列顺序逐个取得所有已发布帧。Frame 对象本身继续由传入的 `std::pmr::memory_resource` 分配,Boost.Lockfree 队列内部节点由 Boost 自己的 allocator 管理。 +`Flow_Refresh_Strategy` 当前使用 `std::pmr::deque>` 保存待渲染帧,不引入 Boost.Lockfree 或其他新增队列依赖。多个 painter 可以并发创建 Frame,入队由 `state_mutex_` 串行化;多个 renderer 调用也安全,但实际消费由 `render_mutex_` 串行化,因此每个 Frame 最多被消费一次。多 producer 下的全局顺序以实际完成入队的线性化顺序为准。 -Flow 的 observer 允许同线程重入 Observer_State;renderer observer 中再次申请同一个 Flow renderer 会立即得到空 lease,不会等待当前 renderer lease 自身释放。 +Flow 的 observer 允许同线程重入 Observer_State;renderer observer 中再次申请同一个 Flow renderer 会立即得到空 lease,不会等待当前 renderer lease 自身释放。完整线程契约和多线程压力测试见 `threading.md`。 ## 实时数据附件所有权 @@ -175,7 +179,6 @@ History_Real_Time_Data> history(memory_resource); 以下内存不属于 Renderive 自己的 allocator-aware 场景域: - Taskflow 4.1.0 内部任务节点、执行队列和 executor 工作资源 -- Boost.Lockfree Flow 队列内部节点 - `std::function` 超出小对象优化后的内部存储 - `std::thread` 的系统线程对象和线程栈 - 异常运行时对象 diff --git a/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.hpp b/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.hpp index 2d094c5..fa0cf67 100644 --- a/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.hpp +++ b/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.hpp @@ -53,6 +53,7 @@ public: static_assert(Timed_Struct_Observer); class Painter_Lease { public: + static constexpr bool thread_affine = true; explicit Painter_Lease(Flow_Refresh_Strategy& strategy); Painter_Lease(const Painter_Lease&) = delete; Painter_Lease& operator=(const Painter_Lease&) = delete; @@ -72,6 +73,7 @@ public: }; class Render_Lease { public: + static constexpr bool thread_affine = true; explicit Render_Lease(Flow_Refresh_Strategy& strategy); Render_Lease(const Render_Lease&) = delete; Render_Lease& operator=(const Render_Lease&) = delete; diff --git a/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.inl b/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.inl index 4a9ecb8..869872d 100644 --- a/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.inl +++ b/Kernel/src/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy.inl @@ -1,8 +1,9 @@ #pragma once template Flow_Refresh_Strategy::Painter_Lease::Painter_Lease(Flow_Refresh_Strategy& strategy) - : strategy_(&strategy), frame_(make_memory_resource_unique(*strategy.memory_resource_)) { + : strategy_(&strategy) { std::lock_guard lock(strategy.state_mutex_); + frame_ = make_memory_resource_unique(*strategy.memory_resource_); frame_->statistics.sequence = ++strategy.next_sequence_; frame_->statistics.real_time_data_update_sequence = strategy.real_time_data_update_sequence_; frame_->statistics.paint_begin_time_ns = strategy.now_ns(); @@ -112,6 +113,7 @@ Flow_Refresh_Strategy::Render_Lease::~Render_Lease strategy_->state_.pending_frame_count = strategy_->frames_.size(); strategy_->update_frame_control_state(Frame_Control_Strategy_Base::invalid_frequency_hz(), 0); observation = {Observation_Event::rendered, frame_->statistics, strategy_->state_, strategy_->last_real_time_data_update_}; + frame_.reset(); } lease_lock_.unlock(); rendering_strategy_ = previous_render_strategy_; diff --git a/Kernel/src/renderive/frame_control/strategy/low_latency/Low_Latency_Strategy.hpp b/Kernel/src/renderive/frame_control/strategy/low_latency/Low_Latency_Strategy.hpp index 8161bcc..bfbe9f9 100644 --- a/Kernel/src/renderive/frame_control/strategy/low_latency/Low_Latency_Strategy.hpp +++ b/Kernel/src/renderive/frame_control/strategy/low_latency/Low_Latency_Strategy.hpp @@ -89,6 +89,7 @@ public: static_assert(Timed_Struct_Observer); class Painter_Lease { public: + static constexpr bool thread_affine = true; explicit Painter_Lease(Low_Latency_Strategy& strategy); Painter_Lease(const Painter_Lease&) = delete; Painter_Lease& operator=(const Painter_Lease&) = delete; @@ -110,6 +111,7 @@ public: }; class Render_Lease { public: + static constexpr bool thread_affine = true; explicit Render_Lease(Low_Latency_Strategy& strategy); Render_Lease(const Render_Lease&) = delete; Render_Lease& operator=(const Render_Lease&) = delete; diff --git a/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.hpp b/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.hpp index 06dd67d..33f9c63 100644 --- a/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.hpp +++ b/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.hpp @@ -50,6 +50,7 @@ public: static_assert(Timed_Struct_Observer); class Painter_Lease { public: + static constexpr bool thread_affine = true; explicit Painter_Lease(Manual_Refresh_Strategy& strategy); Painter_Lease(const Painter_Lease&) = delete; Painter_Lease& operator=(const Painter_Lease&) = delete; @@ -70,6 +71,7 @@ public: }; class Render_Lease { public: + static constexpr bool thread_affine = true; explicit Render_Lease(Manual_Refresh_Strategy& strategy); Render_Lease(const Render_Lease&) = delete; Render_Lease& operator=(const Render_Lease&) = delete; diff --git a/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.inl b/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.inl index bb897b0..5ec5f4b 100644 --- a/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.inl +++ b/Kernel/src/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy.inl @@ -155,7 +155,7 @@ auto Manual_Refresh_Strategy::frame_control_state( } template bool Manual_Refresh_Strategy::refresh() { - std::lock_guard render_lock(render_mutex_); + std::unique_lock render_lock(render_mutex_); Observation observation; bool refreshed{}; { @@ -175,6 +175,7 @@ bool Manual_Refresh_Strategy::refresh() { observation = {Observation_Event::refresh_failed, {}, state_, last_real_time_data_update_}; } } + render_lock.unlock(); observe(observation); return refreshed; } diff --git a/Kernel/src/renderive/real_time_data/History_Real_Time_Data.hpp b/Kernel/src/renderive/real_time_data/History_Real_Time_Data.hpp index 3e7d968..3c61582 100644 --- a/Kernel/src/renderive/real_time_data/History_Real_Time_Data.hpp +++ b/Kernel/src/renderive/real_time_data/History_Real_Time_Data.hpp @@ -41,6 +41,7 @@ public: private: static Container make_container(std::pmr::memory_resource& memory_resource); mutable Mutex mutex_; + std::recursive_mutex mutation_mutex_; Observer observer_; Container values_; std::pmr::vector update_times_; diff --git a/Kernel/src/renderive/real_time_data/History_Real_Time_Data.inl b/Kernel/src/renderive/real_time_data/History_Real_Time_Data.inl index c7ffd06..23d213b 100644 --- a/Kernel/src/renderive/real_time_data/History_Real_Time_Data.inl +++ b/Kernel/src/renderive/real_time_data/History_Real_Time_Data.inl @@ -20,6 +20,7 @@ auto History_Real_Time_Data::make_contai } template void History_Real_Time_Data::update(Value value) { + std::lock_guard mutation_lock(mutation_mutex_); Real_Time_Data_Observation observation; { std::lock_guard lock(mutex_); @@ -34,6 +35,7 @@ void History_Real_Time_Data::update(Valu } template void History_Real_Time_Data::clear() { + std::lock_guard mutation_lock(mutation_mutex_); Real_Time_Data_Observation observation; { std::lock_guard lock(mutex_); @@ -72,6 +74,7 @@ Real_Time_Data_Retention History_Real_Time_Data std::size_t History_Real_Time_Data::discard_before_time_ns(std::uint64_t time_ns) { + std::lock_guard mutation_lock(mutation_mutex_); Real_Time_Data_Observation observation; std::size_t discarded{}; { diff --git a/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.hpp b/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.hpp index 582b1e3..8d0fcaa 100644 --- a/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.hpp +++ b/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.hpp @@ -34,6 +34,7 @@ public: } private: mutable Mutex mutex_; + std::recursive_mutex mutation_mutex_; Observer observer_; std::optional value_; std::uint64_t revision_{}; diff --git a/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.inl b/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.inl index 7046b30..c8e47ff 100644 --- a/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.inl +++ b/Kernel/src/renderive/real_time_data/Latest_Real_Time_Data.inl @@ -3,6 +3,7 @@ template Latest_Real_Time_Data::Latest_Real_Time_Data(Observer observer) : observer_(std::move(observer)) {} template void Latest_Real_Time_Data::update(Value value) { + std::lock_guard mutation_lock(mutation_mutex_); Real_Time_Data_Observation observation; { std::lock_guard lock(mutex_); diff --git a/Kernel/src/renderive/scene/Scene2D_Context.hpp b/Kernel/src/renderive/scene/Scene2D_Context.hpp index 36d39c5..1a8a86d 100644 --- a/Kernel/src/renderive/scene/Scene2D_Context.hpp +++ b/Kernel/src/renderive/scene/Scene2D_Context.hpp @@ -87,7 +87,11 @@ protected: } void on_renderable_detached(Renderable_Base& renderable) override { promote_children(Scene_Base::layer_node(renderable)); - promote_children(Scene_Base::dependency_node(renderable)); + auto& dependency = Scene_Base::dependency_node(renderable); + for (std::size_t index = 0; index < dependency.child_count(); ++index) { + dependency.child(index).owner()->invalidate_cache(); + } + promote_children(dependency); color_caches_.erase(&renderable); } void prepare_render_task(Render_Task& task, const Renderable_List& renderables) override { diff --git a/Kernel/src/renderive/scene/base/Scene_Base.cpp b/Kernel/src/renderive/scene/base/Scene_Base.cpp index 9b0a87a..47b3497 100644 --- a/Kernel/src/renderive/scene/base/Scene_Base.cpp +++ b/Kernel/src/renderive/scene/base/Scene_Base.cpp @@ -15,6 +15,15 @@ public: return executor; } }; +class Scene_Base::Render_Execution_Scope { +public: + explicit Render_Execution_Scope(Scene_Base& scene) noexcept : previous_(std::exchange(active_execution_scene_, &scene)) {} + ~Render_Execution_Scope() { + active_execution_scene_ = previous_; + } +private: + Scene_Base* previous_{}; +}; Scene_Base::Scene_Base() : Scene_Base(*std::pmr::get_default_resource()) {} Scene_Base::Scene_Base(std::pmr::memory_resource& upstream_memory_resource) : scene_lifetime_(std::make_shared(*this)), memory_domain_(std::allocate_shared(std::pmr::polymorphic_allocator(&upstream_memory_resource), upstream_memory_resource)), renderable_states_{Renderable_List(&memory_domain_->resource()), Renderable_List(&memory_domain_->resource())}, render_renderables_(&renderable_states_[0]), cache_renderables_(&renderable_states_[1]), task_(memory_domain_->resource()) { @@ -26,7 +35,19 @@ Scene_Base::~Scene_Base() { shutdown(); } void Scene_Base::render() { + if (active_submitted_observer_scene_ == this) { + std::lock_guard lock(task_mutex_); + ++deferred_render_count_; + return; + } auto task_lock = lock_render_idle(); + if (pending_exception_observed_) { + pending_exception_ = nullptr; + pending_exception_observed_ = false; + } + if (!pending_exception_ && current_completion_ && current_completion_->completed && current_completion_->exception && !current_completion_->observed) { + pending_exception_ = current_completion_->exception; + } Render_Task task(memory_resource()); { std::lock_guard renderable_lock(renderable_mutex_); @@ -41,14 +62,25 @@ void Scene_Base::render() { task.render_sequence = ++render_sequence_; task.completion = std::make_shared(); current_completion_ = task.completion; - observe_scene({Observation_Event::render_submitted, observer_now_ns(), task.render_sequence, task.scene_state_revision, task.render_order.size()}); + const Observation submitted_observation{Observation_Event::render_submitted, observer_now_ns(), task.render_sequence, task.scene_state_revision, task.render_order.size()}; task_ = std::move(task); task_pending_ = true; task_lock.unlock(); + Scene_Base* previous_observer_scene = std::exchange(active_submitted_observer_scene_, this); + observe_scene(submitted_observation); + active_submitted_observer_scene_ = previous_observer_scene; task_ready_.notify_one(); + std::size_t deferred_render_count{}; + { + std::lock_guard lock(task_mutex_); + deferred_render_count = std::exchange(deferred_render_count_, 0); + } + for (std::size_t index = 0; index < deferred_render_count; ++index) { + render(); + } } void Scene_Base::wait_for_render() { - if (is_render_worker_thread()) { + if (is_render_worker_thread() || active_submitted_observer_scene_ == this) { return; } std::unique_lock lock(task_mutex_); @@ -59,7 +91,17 @@ void Scene_Base::wait_for_render() { render_completed_.wait(lock, [&completion] { return completion->completed; }); - const auto exception = completion->exception; + std::exception_ptr exception; + if (pending_exception_) { + exception = pending_exception_; + pending_exception_observed_ = true; + if (!completion->exception) { + completion->observed = true; + } + } else { + exception = completion->exception; + completion->observed = true; + } lock.unlock(); if (exception) { std::rethrow_exception(exception); @@ -218,6 +260,12 @@ std::uint64_t Scene_Base::observer_now_ns() const noexcept { return 0; } std::unique_lock Scene_Base::lock_render_idle() { + if (is_render_worker_thread()) { + throw std::logic_error("scene control mutation is not allowed during render execution"); + } + if (active_submitted_observer_scene_ == this) { + throw std::logic_error("scene control mutation is not allowed during render submission observation"); + } std::unique_lock lock(task_mutex_); render_completed_.wait(lock, [this] { return !task_pending_ && !rendering_; @@ -225,7 +273,7 @@ std::unique_lock Scene_Base::lock_render_idle() { return lock; } bool Scene_Base::is_render_worker_thread() const noexcept { - return worker_.joinable() && std::this_thread::get_id() == worker_.get_id(); + return active_execution_scene_ == this || (worker_.joinable() && std::this_thread::get_id() == worker_.get_id()); } void Scene_Base::shutdown() noexcept { { @@ -323,7 +371,8 @@ void Scene_Base::execute_taskflow(Render_Task& task) { std::pmr::vector internal_tasks(&scratch_resource); internal_tasks.reserve(nodes.size()); for (const auto& node : nodes) { - internal_tasks.push_back(module.graph.emplace([function = node.function, context] { + internal_tasks.push_back(module.graph.emplace([this, function = node.function, context] { + Render_Execution_Scope scope(*this); function(context); }).name(std::string(node.name))); } @@ -359,6 +408,7 @@ void Scene_Base::execute_taskflow(Render_Task& task) { } const Scene_Render_Context final_context{this, task.frame_control_state, task.scene_state_revision, task.render_sequence, nullptr, nullptr}; tf::Task final_task = scene_graph.emplace([this, final_context, &task] { + Render_Execution_Scope scope(*this); generate_final_color_cache(final_context, task.display_order); }).name("final_color_cache"); if (modules.empty()) { diff --git a/Kernel/src/renderive/scene/base/Scene_Base.hpp b/Kernel/src/renderive/scene/base/Scene_Base.hpp index 7e0a9d5..c96b2bc 100644 --- a/Kernel/src/renderive/scene/base/Scene_Base.hpp +++ b/Kernel/src/renderive/scene/base/Scene_Base.hpp @@ -100,15 +100,19 @@ protected: private: friend class Renderable_Base; class Execution_Context; + class Render_Execution_Scope; struct Render_Completion { std::exception_ptr exception; bool completed{}; + bool observed{}; }; void render_loop(); void execute_taskflow(Render_Task& task); void validate_renderable_scene(const Renderable_Base& renderable) const; bool is_renderable_attached_locked(const Renderable_Base& renderable) const; void validate_renderable_attached_locked(const Renderable_Base& renderable) const; + inline static thread_local Scene_Base* active_execution_scene_{}; + inline static thread_local Scene_Base* active_submitted_observer_scene_{}; std::shared_ptr scene_lifetime_; std::shared_ptr memory_domain_; std::array renderable_states_; @@ -121,6 +125,9 @@ private: std::thread worker_; Render_Task task_; std::shared_ptr current_completion_; + std::exception_ptr pending_exception_; + bool pending_exception_observed_{}; + std::size_t deferred_render_count_{}; std::uint64_t render_sequence_{}; bool task_pending_{}; bool rendering_{}; diff --git a/Kernel/src/renderive/scene/base/Scene_Lifetime.hpp b/Kernel/src/renderive/scene/base/Scene_Lifetime.hpp index b655ec8..10cc3ae 100644 --- a/Kernel/src/renderive/scene/base/Scene_Lifetime.hpp +++ b/Kernel/src/renderive/scene/base/Scene_Lifetime.hpp @@ -1,4 +1,6 @@ #pragma once +#include +#include #include #include class Scene_Base; @@ -6,8 +8,22 @@ class Scene_Lifetime { public: class Lease { public: - Lease(Lease&&) noexcept = default; - Lease& operator=(Lease&&) noexcept = default; + Lease() = default; + Lease(const Lease&) = delete; + Lease& operator=(const Lease&) = delete; + Lease(Lease&& other) noexcept : owner_(std::exchange(other.owner_, nullptr)), scene_(std::exchange(other.scene_, nullptr)) {} + Lease& operator=(Lease&& other) noexcept { + if (this == &other) { + return *this; + } + release(); + owner_ = std::exchange(other.owner_, nullptr); + scene_ = std::exchange(other.scene_, nullptr); + return *this; + } + ~Lease() { + release(); + } explicit operator bool() const noexcept { return scene_ != nullptr; } @@ -16,20 +32,44 @@ public: } private: friend class Scene_Lifetime; - Lease(std::unique_lock lock, Scene_Base* scene) noexcept : lock_(std::move(lock)), scene_(scene) {} - std::unique_lock lock_; + Lease(Scene_Lifetime* owner, Scene_Base* scene) noexcept : owner_(owner), scene_(scene) {} + void release() noexcept { + if (!owner_) { + return; + } + owner_->release(); + owner_ = nullptr; + scene_ = nullptr; + } + Scene_Lifetime* owner_{}; Scene_Base* scene_{}; }; explicit Scene_Lifetime(Scene_Base& scene) noexcept : scene_(&scene) {} - Lease acquire() const noexcept { - std::unique_lock lock(mutex_); - return Lease(std::move(lock), scene_); + Lease acquire() noexcept { + std::lock_guard lock(mutex_); + if (!scene_) { + return {}; + } + ++lease_count_; + return Lease(this, scene_); } void invalidate() noexcept { - std::lock_guard lock(mutex_); + std::unique_lock lock(mutex_); scene_ = nullptr; + condition_.wait(lock, [this] { + return lease_count_ == 0; + }); } private: + void release() noexcept { + std::lock_guard lock(mutex_); + --lease_count_; + if (lease_count_ == 0) { + condition_.notify_all(); + } + } mutable std::mutex mutex_; + std::condition_variable condition_; Scene_Base* scene_{}; + std::size_t lease_count_{}; }; diff --git a/Kernel/tests/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy_Test.cpp b/Kernel/tests/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy_Test.cpp index 5ebeb67..544dfe7 100644 --- a/Kernel/tests/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy_Test.cpp +++ b/Kernel/tests/renderive/frame_control/strategy/flow/Flow_Refresh_Strategy_Test.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -180,3 +181,42 @@ TEST(flow_refresh_strategy_test, render_observer_reentry_does_not_deadlock_rende ASSERT_TRUE(second); EXPECT_EQ(second->value, 2); } +TEST(flow_refresh_strategy_test, concurrent_producers_and_consumer_support_unsynchronized_pmr_resource) { + std::pmr::unsynchronized_pool_resource resource; + Flow_Refresh_Test_Strategy strategy(resource); + constexpr int producer_count = 8; + constexpr int frames_per_producer = 200; + constexpr int frame_count = producer_count * frames_per_producer; + std::atomic producers_done{}; + std::atomic consumed{}; + std::thread consumer([&] { + while (consumed.load(std::memory_order_acquire) < frame_count) { + auto frame = strategy.acquire_renderer(); + if (frame) { + 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([&strategy, &producers_done, producer] { + for (int index = 0; index < frames_per_producer; ++index) { + auto frame = strategy.acquire_painter(); + frame->value = producer * frames_per_producer + index; + } + producers_done.fetch_add(1, std::memory_order_release); + }); + } + for (auto& producer : producers) { + producer.join(); + } + consumer.join(); + EXPECT_EQ(consumed.load(std::memory_order_acquire), frame_count); + EXPECT_EQ(strategy.pending_frame_count(), 0); +} diff --git a/Kernel/tests/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy_Test.cpp b/Kernel/tests/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy_Test.cpp index d6c1f25..0be7fb1 100644 --- a/Kernel/tests/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy_Test.cpp +++ b/Kernel/tests/renderive/frame_control/strategy/manual/Manual_Refresh_Strategy_Test.cpp @@ -1,4 +1,8 @@ #include +#include +#include +#include +#include #include "renderive/frame_control/Frame_Control.hpp" struct Manual_Refresh_Test_Frame { int value{}; @@ -44,3 +48,39 @@ TEST(manual_refresh_strategy_test, reports_failed_refresh_without_pending_frame) EXPECT_FALSE(strategy.refresh()); EXPECT_EQ(strategy.state().failed_refresh_count, 1); } +struct Manual_Refresh_Reentrant_Observer_Data { + std::function callback; +}; +struct Manual_Refresh_Reentrant_Observer { + static constexpr bool enabled = true; + std::shared_ptr data{std::make_shared()}; + template + void observe(const Observation& observation) noexcept { + if (data->callback) { + data->callback(static_cast(observation.event)); + } + } +}; +using Manual_Refresh_Reentrant_Observer_State = Observer_State; +using Manual_Refresh_Reentrant_Strategy = Manual_Refresh_Strategy; +TEST(manual_refresh_strategy_test, refresh_observer_can_acquire_renderer_without_deadlock) { + Manual_Refresh_Reentrant_Observer observer; + auto data = observer.data; + Manual_Refresh_Reentrant_Strategy strategy{Manual_Refresh_Reentrant_Observer_State(observer)}; + { + auto frame = strategy.acquire_painter(); + frame->value = 42; + } + std::atomic observed_value{}; + data->callback = [&](int event) { + if (event != static_cast(Manual_Refresh_Reentrant_Strategy::Observation_Event::refresh_succeeded)) { + return; + } + auto frame = strategy.acquire_renderer(); + if (frame) { + observed_value.store(frame->value, std::memory_order_release); + } + }; + EXPECT_TRUE(strategy.refresh()); + EXPECT_EQ(observed_value.load(std::memory_order_acquire), 42); +} diff --git a/Kernel/tests/renderive/real_time_data/Real_Time_Data_Test.cpp b/Kernel/tests/renderive/real_time_data/Real_Time_Data_Test.cpp index 5e2e0c7..34e5cec 100644 --- a/Kernel/tests/renderive/real_time_data/Real_Time_Data_Test.cpp +++ b/Kernel/tests/renderive/real_time_data/Real_Time_Data_Test.cpp @@ -1,6 +1,7 @@ #include #include #include +#include #include #include #include @@ -189,3 +190,33 @@ TEST(real_time_data_attachment_test, attachment_owns_real_time_data_sources) { renderable.reset(); EXPECT_TRUE(weak.expired()); } +struct Real_Time_Data_Reentrant_Frame_Observer_Data { + std::function callback; +}; +struct Real_Time_Data_Reentrant_Frame_Observer { + static constexpr bool enabled = true; + std::shared_ptr data{std::make_shared()}; + template + void observe(const Observation&) noexcept { + if (data->callback) { + data->callback(); + } + } +}; +using Real_Time_Data_Reentrant_Frame_Observer_State = Observer_State; +using Real_Time_Data_Reentrant_Strategy = Low_Latency_Strategy; +TEST(real_time_data_attachment_test, frame_strategy_observer_can_reacquire_scene_lifetime_without_deadlock) { + Real_Time_Data_Reentrant_Frame_Observer observer; + auto observer_data = observer.data; + Real_Time_Data_Reentrant_Frame_Observer_State observer_state(observer, Real_Time_Data_Test_Time_Source{}); + Real_Time_Data_Reentrant_Strategy::Configuration configuration{60.0}; + Scene2D_Context scene(std::move(observer_state), configuration); + auto latest = std::make_shared(); + auto renderable = std::make_shared(With_Real_Time_Data(latest), scene); + std::atomic reacquired{}; + observer_data->callback = [&] { + reacquired.store(&renderable->scene() == &scene, std::memory_order_release); + }; + latest->update(1); + EXPECT_TRUE(reacquired.load(std::memory_order_acquire)); +} diff --git a/Kernel/tests/renderive/scene/Scene2D_Context_Test.cpp b/Kernel/tests/renderive/scene/Scene2D_Context_Test.cpp index d8555c3..6cdbabe 100644 --- a/Kernel/tests/renderive/scene/Scene2D_Context_Test.cpp +++ b/Kernel/tests/renderive/scene/Scene2D_Context_Test.cpp @@ -129,3 +129,31 @@ TEST(scene2d_context_test, final_color_cache_callback_can_reenter_scene_control_ }); EXPECT_FALSE(renderable->configuration().cache_enabled); } +TEST(scene2d_context_test, detach_dependency_parent_invalidates_promoted_cached_child) { + Scene2D_Context<> scene; + auto grandparent = std::make_shared(scene, 1); + auto parent = std::make_shared(scene, 2); + auto child = std::make_shared(scene, 3); + scene.attach_renderable(grandparent); + scene.attach_renderable(parent); + scene.attach_renderable(child); + scene.set_dependency_parent(*parent, grandparent.get()); + scene.set_dependency_parent(*child, parent.get()); + scene.render(); + scene.wait_for_render(); + EXPECT_EQ(grandparent->render_count, 1); + EXPECT_EQ(parent->render_count, 1); + EXPECT_EQ(child->render_count, 1); + scene.detach_renderable(*parent); + scene.render(); + scene.wait_for_render(); + EXPECT_EQ(grandparent->render_count, 1); + EXPECT_EQ(parent->render_count, 1); + EXPECT_EQ(child->render_count, 2); + const auto topology = scene.topology_snapshot(); + for (const auto& relation : topology.dependency) { + if (relation.child.get() == child.get()) { + EXPECT_EQ(relation.parent.get(), grandparent.get()); + } + } +} diff --git a/Kernel/tests/renderive/scene/Scene_State_Observer_Test.cpp b/Kernel/tests/renderive/scene/Scene_State_Observer_Test.cpp index 62b0789..4a5e15d 100644 --- a/Kernel/tests/renderive/scene/Scene_State_Observer_Test.cpp +++ b/Kernel/tests/renderive/scene/Scene_State_Observer_Test.cpp @@ -65,3 +65,21 @@ TEST(scene_state_observer_test, render_worker_observer_can_wait_for_current_rend scene.wait_for_render(); EXPECT_TRUE(reentered.load(std::memory_order_acquire)); } +TEST(scene_state_observer_test, render_submitted_observer_can_submit_next_render_without_deadlock) { + Scene_Reentrant_Observer recorder; + auto data = recorder.data; + using Scene = Scene2D_Context, Recording_Color_Cache, Scene_State_Observer_Test_State, Scene_State_Observer_Test_State_Observer, Scene_Reentrant_Observer_State>; + Scene scene(With_Observer(Scene_State_Observer_Test_State_Observer{}), With_Observer(Scene_Reentrant_Observer_State(recorder, Scene_State_Observer_Test_Time_Source{}))); + std::atomic submitted{}; + data->callback = [&](Scene_Base::Observation_Event event) { + if (event != Scene_Base::Observation_Event::render_submitted) { + return; + } + if (submitted.fetch_add(1, std::memory_order_acq_rel) == 0) { + scene.render(); + } + }; + scene.render(); + scene.wait_for_render(); + EXPECT_EQ(submitted.load(std::memory_order_acquire), 2); +} diff --git a/Kernel/tests/renderive/scene/base/Scene_Base_Test.cpp b/Kernel/tests/renderive/scene/base/Scene_Base_Test.cpp index 9fc60b5..81a9acd 100644 --- a/Kernel/tests/renderive/scene/base/Scene_Base_Test.cpp +++ b/Kernel/tests/renderive/scene/base/Scene_Base_Test.cpp @@ -215,3 +215,54 @@ TEST(scene_base_test, all_waiters_receive_the_same_render_failure) { second.join(); EXPECT_EQ(failures.load(), 2); } +struct Scene_Base_Throw_Once_Renderable : Renderable_Base { + explicit Scene_Base_Throw_Once_Renderable(Scene_Base& scene) : Renderable_Base(scene, {.cache_enabled = false}) {} + void render(const Scene_Render_Context&) override { + if (render_count.fetch_add(1, std::memory_order_acq_rel) == 0) { + throw std::runtime_error("first render failed"); + } + } + std::atomic render_count{}; +}; +TEST(scene_base_test, unobserved_render_failure_survives_next_successful_submission) { + Scene2D_Context<> scene; + auto renderable = std::make_shared(scene); + scene.attach_renderable(renderable); + scene.render(); + scene.render(); + EXPECT_THROW(scene.wait_for_render(), std::runtime_error); + EXPECT_EQ(renderable->render_count.load(std::memory_order_acquire), 2); + scene.render(); + EXPECT_NO_THROW(scene.wait_for_render()); + EXPECT_EQ(renderable->render_count.load(std::memory_order_acquire), 3); +} +struct Scene_Base_Render_Execution_Renderable : Renderable_Base { + explicit Scene_Base_Render_Execution_Renderable(Scene_Base& scene) : Renderable_Base(scene, {.cache_enabled = false}) {} + void render(const Scene_Render_Context& context) override { + context.scene->wait_for_render(); + wait_returned.store(true, std::memory_order_release); + try { + context.scene->render(); + } catch (const std::logic_error&) { + nested_render_rejected.store(true, std::memory_order_release); + } + try { + context.scene->set_renderable_configuration(*this, {.cache_enabled = true}); + } catch (const std::logic_error&) { + mutation_rejected.store(true, std::memory_order_release); + } + } + std::atomic wait_returned{}; + std::atomic nested_render_rejected{}; + std::atomic mutation_rejected{}; +}; +TEST(scene_base_test, taskflow_execution_thread_cannot_wait_on_or_mutate_current_scene) { + Scene2D_Context<> scene; + auto renderable = std::make_shared(scene); + scene.attach_renderable(renderable); + scene.render(); + scene.wait_for_render(); + EXPECT_TRUE(renderable->wait_returned.load(std::memory_order_acquire)); + EXPECT_TRUE(renderable->nested_render_rejected.load(std::memory_order_acquire)); + EXPECT_TRUE(renderable->mutation_rejected.load(std::memory_order_acquire)); +} diff --git a/Kernel/tests/renderive/threading/Threading_Contract_Test.cpp b/Kernel/tests/renderive/threading/Threading_Contract_Test.cpp new file mode 100644 index 0000000..3205037 --- /dev/null +++ b/Kernel/tests/renderive/threading/Threading_Contract_Test.cpp @@ -0,0 +1,619 @@ +#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(Scene_Base& scene, std::atomic& render_count, std::atomic& maximum_sequence) + : Renderable_Base(scene, {.cache_enabled = false}), render_count(&render_count), maximum_sequence(&maximum_sequence) {} + void render(const Scene_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.render_sequence && !maximum_sequence->compare_exchange_weak(maximum, context.render_sequence, std::memory_order_acq_rel)) {} + } + std::atomic* render_count; + std::atomic* maximum_sequence; +}; +struct Threading_Test_Task_Graph_Renderable : Renderable_Base { + Threading_Test_Task_Graph_Renderable(Scene_Base& scene, std::atomic& executed) : Renderable_Base(scene, {.cache_enabled = false}), executed(&executed) {} + void build_task_graph(Renderable_Task_Graph& graph) override { + graph.emplace([this](const Scene_Render_Context&) { + executed->fetch_add(1, std::memory_order_relaxed); + }); + } + 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_returns_consistent_snapshots_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 render_reader([&] { + wait_start(start); + while (!writer_done.load(std::memory_order_acquire)) { + if (!valid_frame(strategy.render_use_state().frame)) { + invalid.fetch_add(1, std::memory_order_relaxed); + } + } + }); + 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(); + render_reader.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(); + if (!valid_frame(strategy.render_use_state().frame)) { + invalid.fetch_add(1, std::memory_order_relaxed); + } + } + strategy.acquire_render_state(); + if (!valid_frame(strategy.render_use_state().frame)) { + invalid.fetch_add(1, std::memory_order_relaxed); + } + }); + 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 = std::make_shared(scene, executed, maximum_sequence); + scene.attach_renderable(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_task_graph_rebuild_can_run_concurrently) { + Scene2D_Context<> scene; + constexpr int render_count = 300; + std::atomic executed{}; + auto renderable = std::make_shared(scene, executed); + scene.attach_renderable(renderable); + std::atomic start{}; + std::atomic renderer_done{}; + std::thread renderer([&] { + wait_start(start); + for (int index = 0; index < render_count; ++index) { + scene.render(); + scene.wait_for_render(); + } + renderer_done.store(true, std::memory_order_release); + }); + std::thread rebuilder([&] { + wait_start(start); + while (!renderer_done.load(std::memory_order_acquire)) { + renderable->rebuild_task_graph(); + std::this_thread::yield(); + } + }); + start.store(true, std::memory_order_release); + renderer.join(); + rebuilder.join(); + EXPECT_EQ(executed.load(std::memory_order_acquire), render_count); +} +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(Scene_Base& scene) : Renderable_Base(scene) {} + }; + using Attachment = Attach_Real_Time_Data; + auto data = std::make_shared(); + auto scene = std::make_unique>(); + auto renderable = std::make_shared(With_Real_Time_Data(data), *scene); + 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); +} diff --git a/Kernel/threading.md b/Kernel/threading.md new file mode 100644 index 0000000..9df10e1 --- /dev/null +++ b/Kernel/threading.md @@ -0,0 +1,103 @@ +# Renderive 线程模型 + +## 总原则 + +同一个 Renderive 对象的析构不能和调用方主动发起的普通成员函数并发执行。唯一特意处理的跨生命周期路径是实时数据通知到 Scene:实时数据更新允许与 Scene 析构竞争,`Scene_Lifetime` 会阻止通知进入已经失效的 Scene。 + +Frame lease 是线程绑定的独占操作句柄。内置 `Painter_Lease` 和 `Render_Lease` 都显式声明 `thread_affine = true`。lease 可以在取得它的线程内 move,但不能把仍持有底层互斥锁的 lease 转移到另一线程使用或析构;Frame 内容也不能被多个线程同时操作。 + +用户实现的 Renderable 任务会由 Taskflow 并行执行。不同 Renderable 如果没有依赖关系可以并行,同一个 Renderable 内部没有依赖边的任务也可以并行。用户任务共享的数据必须由用户自己提供同步;Kernel 只保证自己的 Scene、状态、缓存元数据和任务图描述不会产生数据竞争。 + +## Scene + +`render()` 支持多个调用线程同时提交。Scene 不保存多任务提交队列,而是使用 `task_mutex_` 将提交串行化;后一个 `render()` 会等待当前帧结束后再生成下一帧,因此不会覆盖已有提交。 + +`wait_for_render()` 支持多个等待线程。同一次失败渲染的所有 waiter 都读取同一个 `Render_Completion`,不会由第一个 waiter 消耗异常。若失败帧尚未被任何 waiter 观察就提交了下一帧,失败会保存为 pending exception,并由后续 `wait_for_render()` 报告,不会被新的成功 completion 覆盖。 + +`attach_renderable()`、`detach_renderable()`、`set_display_parent()`、`set_dependency_parent()`、`set_renderable_configuration()` 与后台渲染互斥。它们通过 `lock_render_idle()` 等待 Scene 空闲后再修改 Renderable 集合或拓扑。 + +`topology_snapshot()` 和 `renderable_count()` 可以与 Scene 拓扑修改并发调用。拓扑快照持有 `shared_ptr`,所以返回以后即使对应对象从 Scene detach,快照内对象生命周期仍然有效。 + +`with_final_color_cache()` 在普通线程上等待后台渲染结束后读取最终缓存。若从 Scene 自己的 render worker 调用,则直接读取当前 worker 可见的缓存。回调在 Scene 的控制锁范围内执行,不应启动另一个线程并等待该线程重新进入同一个 Scene 控制 API。 + +Scene observer 在触发事件的线程同步执行。`render_submitted` observer 中再次调用 `render()` 会登记为 deferred render,在当前 submitted callback 返回并启动当前帧后顺序提交;在该 callback 内调用 `wait_for_render()` 为 no-op,避免同步 observer 自锁。render worker 和 Taskflow executor worker 中调用 `wait_for_render()` 同样为 no-op;这些 render execution 线程调用 `render()` 或拓扑/配置修改接口会直接抛出 `std::logic_error`,不会等待当前帧形成自锁。 + +## Renderable + +`invalidate_cache()`、`cache_revision()`、`rendered_cache_revision()` 使用原子 revision,可以和后台渲染并发。渲染开始时捕获 revision,结束后只发布该次捕获值,因此渲染过程中发生的 invalidate 不会丢失。 + +`configuration()` 返回值快照;`set_renderable_configuration()` 通过 Scene 串行修改内部配置。 + +`task_graph()` 和 `rebuild_task_graph()` 由 `task_graph_mutex_` 串行化。`task_graph()` 返回不可变 `shared_ptr`,旧图在正在执行的帧释放快照前不会失效。`build_task_graph()` 是构图回调,不应在同一个 Renderable 上再次调用 `task_graph()` 或 `rebuild_task_graph()`。 + +Renderable 可以晚于 Scene 析构以完成自身释放,但 `scene()` 只用于调用方已经保证 Scene 生命周期有效的同步访问。不要让 `scene()` 返回的引用与 Scene 析构竞争。Kernel 内部需要跨 Scene 生命周期访问的实时数据路径使用 `Scene_Lifetime::Lease`,不依赖这个裸引用。 + +## Flow 帧策略 + +当前 `Flow_Refresh_Strategy` 不依赖 Boost.Lockfree,也没有其他新增队列依赖。Frame 使用传入的 `std::pmr::memory_resource` 分配,待渲染帧保存在 `std::pmr::deque` 中。 + +多个 painter 线程可以同时调用 `acquire_painter()`。Frame 分配、队列容器操作和 Frame 回收都纳入 `state_mutex_` 的同一同步边界,因此即使调用方传入 `std::pmr::unsynchronized_pool_resource`,Kernel 也不会并发进入该资源。多个 renderer 调用线程也不会产生数据竞争;`render_mutex_` 保证一次只有一个 renderer lease 真正消费队列,所以每个已入队 Frame 最多消费一次。 + +Flow 的并发语义是线程安全队列访问,不是 lock-free 保证。单线程生产时保持 FIFO;多个生产线程时,全局顺序由实际进入 `state_mutex_` 并完成入队的线性化顺序决定。 + +## Manual 帧策略 + +多个 painter 调用由 `painter_mutex_` 串行化。`refresh()` 和 renderer 共用 `render_mutex_`,所以不会在 renderer 正在读取 render frame 时交换 pending/render 指针。Frame 指针和状态切换同时受 `state_mutex_` 保护。 + +Manual 只保证最新 pending frame 被 refresh;并发 producer 产生的中间帧可以按策略语义被后来的 prepared frame 替换,这不是丢帧 bug。 + +## Low Latency 帧策略 + +多个 painter 调用由 `painter_mutex_` 串行化,多个 renderer 调用由 `render_mutex_` 串行化。paint/cache/render 三个 Frame 指针只在 `state_mutex_` 下交换,因此 painter 写入和 renderer 读取不会落到同一个 Frame 上。 + +`set_frequency_hz()`、实时数据通知、pending frame discard、状态读取都和 Frame 指针切换使用同一个 `state_mutex_`。频率的发布状态另外通过 `Frame_Control_Strategy_Base` 的内部 mutex 保护。 + +## 实时数据 + +`Latest_Real_Time_Data` 的 update/snapshot/state 访问由数据 mutex 保护。`History_Real_Time_Data` 的 update/clear/discard/snapshot/state 同样受数据 mutex 保护。`Scene_Lifetime::Lease` 使用活动 lease 计数保护 Scene 生命周期,不再把 lifetime mutex 持有到 Frame Strategy observer callback 结束,因此 observer 中再次通过 Renderable 获取同一 Scene 不会递归锁死。 + +多线程 mutation 还使用独立 `mutation_mutex_` 将一次 mutation 与对应 observer 通知作为一个顺序单元。这样 revision 1 的通知一定先于 revision 2,不会出现较新的实时数据已经通知 Frame Strategy 后,又被较旧 observation 覆盖 `last_real_time_data_update_` 的情况。`mutation_mutex_` 使用 recursive mutex,使 observer 同线程重入实时数据 mutation 时不会因为通知顺序锁自身死锁。 + +实时数据 attachment 持有数据源 `shared_ptr`。bind/unbind 与并发 update 由 observer binding state 同步;Scene 析构和实时数据 update 可以并发,通知在获取不到有效 `Scene_Lifetime::Lease` 时直接跳过。 + +## 状态策略 + +Double/Triple State 的 set/get/publish/acquire/render snapshot 都在策略 mutex 下完成。对外读取返回值快照,不把内部 buffer 引用暴露到锁外。 + +Triple State 的 render acquire 与前台 publish 可以并发;每次 acquire 得到完整的某一个 published revision,不会观察到半更新 State。 + +## Observer + +默认 `Observer_State` 使用 recursive mutex,将多个线程同时触发的 observer callback 串行化,并允许同一线程的 observer callback 再次进入同一个 Observer_State。 + +Observer callback 是同步调用,不是异步事件队列。回调耗时会直接增加触发线程耗时;用户 observer 自身创建的外部对象生命周期和跨对象锁顺序仍由用户负责。 + +## 多线程测试 + +`tests/renderive/threading/Threading_Contract_Test.cpp` 专门验证线程契约: + +- 8 个线程同时进入 `Observer_State`,确认 observer callback 最大并发数为 1。 +- Flow 使用 4 producer + 4 consumer,验证所有 Frame 恰好消费一次、无重复、无撕裂、pending 计数归零;另用 8 producer + `std::pmr::unsynchronized_pool_resource` 验证 Frame 分配和回收不会并发进入非线程安全 PMR。 +- Manual 使用 4 producer,同时运行 refresher 和 renderer,验证 Frame 指针交换期间不会读到撕裂 Frame。 +- Low Latency 使用 4 producer,同时运行 renderer、frequency configurator 和实时数据 updater,验证三缓冲 Frame 与状态更新并发安全。 +- Latest 实时数据使用 4 updater + snapshot reader,验证 revision、update count 和 value 快照一致。 +- 实时数据 observer 使用 8 updater,验证 observer 收到的 revision 严格按 mutation revision 递增。 +- History 实时数据同时 update、snapshot、discard,验证容器和统计状态一致。 +- Double State 同时 set/publish、读取 cache 和 render snapshot,验证 State 不出现撕裂。 +- Triple State 同时 set/publish、acquire render state 和读取 cache,验证 published revision 快照一致。 +- Scene 使用 4 个 render submitter,验证并发提交不会覆盖任务,最终 render sequence 与实际执行次数一致。 +- `render_submitted` observer 重入 `render()`,验证重入提交会 deferred 而不是和 Scene worker 的 observer 锁形成死锁。 +- Manual `refresh_succeeded` observer 重入 `acquire_renderer()`,验证用户 callback 运行时不再持有 `render_mutex_`。 +- RTD -> Frame Strategy observer 中再次获取 Renderable 的 Scene,验证 `Scene_Lifetime` lease 不跨用户 callback 持有 lifetime mutex。 +- 第一帧失败且未 wait、随后第二帧成功,验证第一帧异常仍由后续 `wait_for_render()` 报告。 +- detach dependency parent 后自动提升 child,验证被提升的缓存 child 会 invalidate 并重新渲染。 +- Renderable task 中调用 `wait_for_render()`,验证当前 Scene execution thread 不等待自身;同线程调用 `render()` 或配置修改会明确拒绝。 +- Built-in Frame lease 明确声明 `thread_affine = true`,测试固定同线程生命周期契约。 +- Scene render 与 `rebuild_task_graph()` 并发执行,验证不可变任务图快照不会在执行期间失效。 +- 实时数据 update 与 Scene 析构竞争,验证 Scene lifetime binding 不产生 UAF。 + +原有测试另外覆盖 render 期间 cache invalidate、任务图构建期间 rebuild、Scene topology 写入与 snapshot 并发、多个 render failure waiter、RTD unbind 与 update 并发等针对性边界。 + +真实 Taskflow 存在时还会运行独立内部任务并行执行测试;内置 dependency-graph executor 只验证 DAG 关系,不伪造并行能力。 + +Linux/GCC 或 Clang 环境建议同时跑 ThreadSanitizer。普通单元测试负责语义和确定性不变量,TSan 负责检测测试覆盖路径上的实际 data race;两者不能互相替代。 diff --git a/Kernel/设计文档.md b/Kernel/设计文档.md deleted file mode 100644 index f013a83..0000000 --- a/Kernel/设计文档.md +++ /dev/null @@ -1,12 +0,0 @@ -renderable 可以被渲染的基础元素 - -状态策略 给类的状态提供线程安全访问更新机制 - -帧策略 提供帧的发布 生成机制 - -plot 上下文对象聚合所有策略 - -使用这个无锁的数据结构 -https://github.com/boostorg/lockfree/tree/boost-1.91.0 - -使用MPSC的队列 之前用的都是错误的 \ No newline at end of file