更新三缓冲

This commit is contained in:
2026-08-24 08:12:20 +08:00
parent db251b5abf
commit d6c1db9d06
26 changed files with 1099 additions and 259 deletions
@@ -0,0 +1,140 @@
#pragma once
#include <concepts>
#include <concurrentqueue-1.0.5/concurrentqueue.h>
#include <cstddef>
#include <functional>
#include <mutex>
#include <new>
#include <shared_mutex>
#include <span>
#include <tuple>
#include <type_traits>
#include <utility>
#include <vector>
namespace double_buffer {
/* Def mechanism declaration. Tag selects one independently exchanged MPMC stream. */
template <typename Tag, typename Value>
struct Mpmc_Triple_Buffer {
using Tag_Type = Tag;
using Value_Type = Value;
};
namespace detail {
/*
* The three buffers are queues, not snapshots. Producers append to collecting while
* the renderer owns rendering and observers own query. advance() rotates roles as:
* completed rendering -> query, collecting -> rendering, expired query -> collecting.
*/
template <typename Declaration>
class Mpmc_Triple_Buffer_Storage {
public:
using Value = typename Declaration::Value_Type;
Mpmc_Triple_Buffer_Storage() = default;
Mpmc_Triple_Buffer_Storage(const Mpmc_Triple_Buffer_Storage&) = delete;
Mpmc_Triple_Buffer_Storage& operator=(const Mpmc_Triple_Buffer_Storage&) = delete;
void submit(Value value) {
std::shared_lock lock(exchange_lock_);
if (!collecting_->enqueue(std::move(value))) throw std::bad_alloc{};
}
void advance() {
std::unique_lock lock(exchange_lock_);
auto* expired_query = query_;
query_ = rendering_;
rendering_ = collecting_;
collecting_ = expired_query;
clear(*collecting_);
}
template <typename Callback>
requires std::invocable<Callback&, std::span<Value>>
decltype(auto) access_rendering(Callback&& callback) {
std::unique_lock lock(exchange_lock_);
return access(*rendering_, std::forward<Callback>(callback));
}
template <typename Callback>
requires std::invocable<Callback&, std::span<const Value>>
decltype(auto) access_query(Callback&& callback) const {
std::unique_lock lock(exchange_lock_);
auto adapter = [&](std::span<Value> values) -> decltype(auto) {
return std::invoke(callback, std::span<const Value>{values});
};
return access(*query_, adapter);
}
private:
using Queue = moodycamel::ConcurrentQueue<Value>;
static void clear(Queue& queue) {
Value value;
while (queue.try_dequeue(value)) {}
}
static void restore(Queue& queue, std::vector<Value>& values) {
for (auto& value : values)
if (!queue.enqueue(std::move(value))) throw std::bad_alloc{};
}
template <typename Callback>
static decltype(auto) access(Queue& queue, Callback&& callback) {
std::vector<Value> values;
values.reserve(queue.size_approx());
Value value;
while (queue.try_dequeue(value)) values.push_back(std::move(value));
if constexpr (std::is_void_v<std::invoke_result_t<Callback&, std::span<Value>>>) {
std::invoke(callback, std::span<Value>{values});
restore(queue, values);
}
else {
static_assert(!std::is_reference_v<std::invoke_result_t<Callback&, std::span<Value>>>,
"triple-buffer access results must not outlive the locked queue view");
auto result = std::invoke(callback, std::span<Value>{values});
restore(queue, values);
return result;
}
}
mutable std::shared_mutex exchange_lock_; /* Serializes role exchange and stable role inspection. */
mutable Queue first_{}; /* First physical MPMC queue; role changes only under exchange_lock_. */
mutable Queue second_{}; /* Second physical MPMC queue; role changes only under exchange_lock_. */
mutable Queue third_{}; /* Third physical MPMC queue; role changes only under exchange_lock_. */
Queue* collecting_{&first_}; /* Concurrent producer destination until the next exchange. */
Queue* rendering_{&second_}; /* Immutable-to-producers queue consumed by the renderer. */
mutable Queue* query_{&third_}; /* Last completed render input exposed to observers. */
};
template <typename Tuple>
class Mpmc_Triple_Buffer_Storage_Set;
template <typename... Declarations>
class Mpmc_Triple_Buffer_Storage_Set<std::tuple<Declarations...>> {
template <typename Tag, typename First, typename... Rest>
static consteval std::size_t index_of() {
if constexpr (std::same_as<Tag, typename First::Tag_Type>) return 0;
else return 1 + index_of<Tag, Rest...>();
}
public:
template <typename Tag>
auto& get() noexcept {
constexpr std::size_t index = index_of<Tag, Declarations...>();
return std::get<index>(storage_);
}
template <typename Tag>
const auto& get() const noexcept {
return const_cast<Mpmc_Triple_Buffer_Storage_Set*>(this)->template get<Tag>();
}
private:
std::tuple<Mpmc_Triple_Buffer_Storage<Declarations>...> storage_;
};
}
}
+24 -1
View File
@@ -1,4 +1,5 @@
#pragma once
#include "Mpmc_Triple_Buffer.hpp"
#include <algorithm>
#include <atomic>
#include <cstddef>
@@ -212,6 +213,10 @@ struct Is_Dependency_Graph_Type : std::false_type {};
template <typename Tag, typename Object>
struct Is_Dependency_Graph_Type<Dependency_Graph_Type<Tag, Object>> : std::true_type {};
template <typename Value>
struct Is_Mpmc_Triple_Buffer : std::false_type {};
template <typename Tag, typename Value>
struct Is_Mpmc_Triple_Buffer<Mpmc_Triple_Buffer<Tag, Value>> : std::true_type {};
template <typename Value>
struct Is_State_Type : std::false_type {};
template <typename Tag, typename Prev>
struct Is_State_Type<State_Type<Tag, Prev>> : std::true_type {};
@@ -224,7 +229,9 @@ concept Buffer_Type = Is_Tagged_Buffer<Value>::value;
template <typename Value>
concept Dependency_Graph_Mechanism = Is_Dependency_Graph_Type<Value>::value;
template <typename Value>
concept Mechanism_Type = Buffer_Type<Value> || Dependency_Graph_Mechanism<Value>;
concept Mpmc_Triple_Buffer_Mechanism = Is_Mpmc_Triple_Buffer<Value>::value;
template <typename Value>
concept Mechanism_Type = Buffer_Type<Value> || Dependency_Graph_Mechanism<Value> || Mpmc_Triple_Buffer_Mechanism<Value>;
template <typename Value>
concept Tagged_State = Prop_State<Value> && requires {
typename Value::Tag_Type;
@@ -298,6 +305,8 @@ template <typename T>
using Mechanism_Buffer_Tuple = std::conditional_t<Buffer_Type<T>, std::tuple<T>, std::tuple<>>;
template <typename T>
using Mechanism_Dependency_Graph_Tuple = std::conditional_t<Dependency_Graph_Mechanism<T>, std::tuple<T>, std::tuple<>>;
template <typename T>
using Mechanism_Mpmc_Triple_Buffer_Tuple = std::conditional_t<Mpmc_Triple_Buffer_Mechanism<T>, std::tuple<T>, std::tuple<>>;
template <typename Tag, typename Tuple>
struct Has_Tag : std::false_type {};
template <typename Tag, typename... Types>
@@ -328,6 +337,18 @@ template <typename Tuple>
concept Buffer_List = Is_Buffer_List<Tuple>::value;
template <typename Tuple>
concept Dependency_Graph_List = Is_Dependency_Graph_List<Tuple>::value;
template <typename Tuple>
struct Is_Mpmc_Triple_Buffer_List : std::false_type {};
template <typename... Declarations>
struct Is_Mpmc_Triple_Buffer_List<std::tuple<Declarations...>> : Tagged_List_Check<
(Mpmc_Triple_Buffer_Mechanism<Declarations> && ...), Declarations...> {};
template <typename Tuple>
concept Mpmc_Triple_Buffer_List = Is_Mpmc_Triple_Buffer_List<Tuple>::value;
template <typename Tag, typename Tuple>
concept Mpmc_Triple_Buffer_Tag_In = requires {
requires Mpmc_Triple_Buffer_List<Tuple>;
requires Has_Tag<Tag, Tuple>::value;
};
template <typename Tag, typename Tuple>
concept Buffer_Tag_In = requires {
requires Buffer_List<Tuple>;
@@ -655,6 +676,7 @@ concept Object_Core = requires {
requires detail::State_Chain_Matches<typename T::State, typename T::States>;
requires detail::Buffer_List<typename T::Buffers>;
requires detail::Dependency_Graph_List<typename T::Dependency_Graph_Types>;
requires detail::Mpmc_Triple_Buffer_List<typename T::Mpmc_Triple_Buffers>;
};
namespace detail {
template <typename Tuple>
@@ -678,6 +700,7 @@ struct Root {
};
using Buffers = std::tuple<>;
using Dependency_Graph_Types = std::tuple<>;
using Mpmc_Triple_Buffers = std::tuple<>;
using Base_Tag = Root;
using States = std::tuple<State_Type<Base_Tag>>;
struct Prop : Prop_Type<Base_Tag> {
+27
View File
@@ -35,6 +35,11 @@ using Impl_Dependency_Graph_Types = decltype(std::tuple_cat(
std::declval<typename Base::Dependency_Graph_Types>(),
std::declval<Mechanism_Dependency_Graph_Tuple<Local_Mechanisms>>()...
));
template <typename Base, typename... Local_Mechanisms>
using Impl_Mpmc_Triple_Buffers = decltype(std::tuple_cat(
std::declval<typename Base::Mpmc_Triple_Buffers>(),
std::declval<Mechanism_Mpmc_Triple_Buffer_Tuple<Local_Mechanisms>>()...
));
template <typename Self, typename Base>
using Impl_States = decltype(std::tuple_cat(
std::declval<typename Base::States>(),
@@ -44,6 +49,7 @@ template <typename Self, typename Base, typename... Local_Mechanisms>
concept Impl_Mechanisms = Object_Root<Base> && (Mechanism_Type<Local_Mechanisms> && ...) && requires {
requires Buffer_List<Impl_Buffers<Base, Local_Mechanisms...>>;
requires Dependency_Graph_List<Impl_Dependency_Graph_Types<Base, Local_Mechanisms...>>;
requires Mpmc_Triple_Buffer_List<Impl_Mpmc_Triple_Buffers<Base, Local_Mechanisms...>>;
requires State_List<Impl_States<Self, Base>>;
};
template <typename Callback, typename Private, typename Object_T>
@@ -65,6 +71,7 @@ struct Def : Base {
using Base_Private = typename Base::Private;
using Buffers = detail::Impl_Buffers<Base, Local_Mechanisms...>;
using Dependency_Graph_Types = detail::Impl_Dependency_Graph_Types<Base, Local_Mechanisms...>;
using Mpmc_Triple_Buffers = detail::Impl_Mpmc_Triple_Buffers<Base, Local_Mechanisms...>;
using States = detail::Impl_States<Self, Base>;
struct Private : Base_Private {
using Tag_Type = Base_Tag;
@@ -118,6 +125,7 @@ struct Impl : Obj {
using State = typename Obj::State;
using Buffers = typename Obj::Buffers;
using Dependency_Graph_Types = typename Obj::Dependency_Graph_Types;
using Mpmc_Triple_Buffers = typename Obj::Mpmc_Triple_Buffers;
using States = typename Obj::States;
/* CRTP 最终 Builder:所有基类流式接口都返回本类型,确保 build() 分派到最派生业务 Builder。 */
struct Builder : Obj::template Builder<Impl> {
@@ -131,6 +139,7 @@ struct Impl : Obj {
detail::State_Callback_Storage<typename Obj::State> state_callbacks;
detail::Buffer_Storage<Buffers> buffer_storage;
detail::Dependency_Graph_Storage<Dependency_Graph_Types> dependency_graph_storage;
detail::Mpmc_Triple_Buffer_Storage_Set<Mpmc_Triple_Buffers> mpmc_triple_buffer_storage;
explicit Private(std::pmr::memory_resource* resource) : Publish_Double_Buffer<Prop>(std::allocator_arg, Allocator{resource}),
state(std::allocator_arg, Allocator{resource}),
buffer_storage(std::allocator_arg, Allocator{resource}),
@@ -244,6 +253,24 @@ public:
this->emit_dependency_source(detail::dependency_id<detail::Buffer_Dependency_Key<Tag>>());
return *data().buffer_storage.template get<Tag>().pending;
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag, typename Input>
void submit_stream(Input&& input) {
data().mpmc_triple_buffer_storage.template get<Tag>().submit(std::forward<Input>(input));
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag>
void exchange_stream() {
data().mpmc_triple_buffer_storage.template get<Tag>().advance();
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag, typename Callback>
decltype(auto) access_rendering_stream(Callback&& callback) {
return data().mpmc_triple_buffer_storage.template get<Tag>().access_rendering(
std::forward<Callback>(callback));
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag, typename Callback>
decltype(auto) access_query_stream(Callback&& callback) const {
return data().mpmc_triple_buffer_storage.template get<Tag>().access_query(
std::forward<Callback>(callback));
}
template <detail::Buffer_Tag_In<Buffers> Tag>
const auto& current_buffer() const {
return *data().buffer_storage.template get<Tag>().current;
+1
View File
@@ -21,6 +21,7 @@ using double_buffer::Prop_Access;
using double_buffer::Pmr;
using double_buffer::Root;
using double_buffer::Tagged_Buffer;
using double_buffer::Mpmc_Triple_Buffer;
using double_buffer::Attached;
using double_buffer::Dependency_Graph_Type;
using double_buffer::Dependency_Graph;
+3 -22
View File
@@ -7,32 +7,13 @@ Scene::Private::~Private() {
}
void Scene::Private::push_event(Event_Pointer event) {
if (!event) throw std::invalid_argument("scene event ownership must not be empty");
std::lock_guard lock(runtime->event_mutex);
runtime->events.pending->push_back(std::move(event));
runtime->submit_event(runtime->object, std::move(event));
}
Scene::Event_Batch Scene::Private::take_events(std::uint64_t frame_sequence) {
Event_Batch result{runtime->events.current->get_allocator()};
{
std::lock_guard lock(runtime->event_mutex);
runtime->events.advance();
runtime->events.pending->clear();
result = std::move(*runtime->events.current);
}
for (const auto& event : result) event->mark_dispatch_started(frame_sequence);
if (!result.empty()) {
std::lock_guard lock(runtime->report_mutex);
runtime->reports.pending->insert(runtime->reports.pending->end(),
result.begin(), result.end());
}
return result;
return runtime->take_events(runtime->object, runtime->resource, frame_sequence);
}
Scene::Event_Report_Batch Scene::Private::take_event_reports() {
Event_Report_Batch result{runtime->reports.current->get_allocator()};
std::lock_guard lock(runtime->report_mutex);
runtime->reports.advance();
runtime->reports.pending->clear();
result = std::move(*runtime->reports.current);
return result;
return runtime->take_reports(runtime->object, runtime->resource);
}
void Scene::dispatch_event(Event_Pointer event) {
static_cast<Private&>(*d).push_event(std::move(event));
+4 -1
View File
@@ -5,12 +5,15 @@
#include <memory_resource>
#include <vector>
namespace aethera {
struct Scene_Event_Stream_Tag {};
/* Scene 状态标签,用于访问和订阅 Scene::State。 */
/*
* Scene 只汇总跨渲染后端共有的 Prepare 数据依赖图并构建 Taskflow。
* 用户最终通过 Impl<Scene> 创建可使用实例;编辑 Dependency_Graph 后调用 advance() 提交结构变化,再调用 process(...) 执行当前场景。
*/
struct Scene : Def<Scene, Root, Dependency_Graph_Type<Prepare_Data_Tag, Renderable>> {
struct Scene : Def<Scene, Root,
Dependency_Graph_Type<Prepare_Data_Tag, Renderable>,
Mpmc_Triple_Buffer<Scene_Event_Stream_Tag, std::shared_ptr<Event>>> {
using Event_Pointer = std::shared_ptr<Event>;
using Event_Batch = std::pmr::vector<Event_Pointer>;
using Event_Report_Pointer = std::shared_ptr<const Event>;
+38 -8
View File
@@ -1,6 +1,6 @@
#pragma once
#include <algorithm>
#include <mutex>
#include <span>
#include <unordered_map>
#include <unordered_set>
#include <vector>
@@ -32,18 +32,48 @@ struct Scene::Private : Prev_Private {
};
struct Scene::Private::Runtime {
std::unique_ptr<tf::Taskflow> taskflow; /* 当前已构建的总 Taskflow;为空表示尚未构建。 */
std::mutex event_mutex;
double_buffer::Double_Buffer<Event_Batch> events;
std::mutex report_mutex;
double_buffer::Double_Buffer<Event_Report_Batch> reports;
explicit Runtime(std::pmr::memory_resource* resource)
: events(std::allocator_arg, typename decltype(events)::allocator_type{resource}),
reports(std::allocator_arg, typename decltype(reports)::allocator_type{resource}) {}
using Submit_Event_Call = void (*)(Root*, Event_Pointer);
using Take_Events_Call = Event_Batch (*)(Root*, std::pmr::memory_resource*, std::uint64_t);
using Take_Reports_Call = Event_Report_Batch (*)(const Root*, std::pmr::memory_resource*);
std::pmr::memory_resource* resource;
Root* object{};
Submit_Event_Call submit_event{};
Take_Events_Call take_events{};
Take_Reports_Call take_reports{};
explicit Runtime(std::pmr::memory_resource* resource) : resource(resource) {}
};
template <Attached Object>
void Scene::Private::bind_private_crtp(Object* object) {
Prev_Private::bind_private_crtp(object);
runtime = std::make_unique<Runtime>(object->memory_resource());
runtime->object = object;
runtime->submit_event = [](Root* root, Event_Pointer event) {
static_cast<Object*>(root)->template submit_stream<Scene_Event_Stream_Tag>(std::move(event));
};
runtime->take_events = [](Root* root, std::pmr::memory_resource* resource,
std::uint64_t frame_sequence) {
auto* scene = static_cast<Object*>(root);
scene->template exchange_stream<Scene_Event_Stream_Tag>();
return scene->template access_rendering_stream<Scene_Event_Stream_Tag>(
[&](std::span<Event_Pointer> events) {
Event_Batch result{resource};
result.reserve(events.size());
for (const auto& event : events) {
event->mark_dispatch_started(frame_sequence);
result.push_back(event);
}
return result;
});
};
runtime->take_reports = [](const Root* root, std::pmr::memory_resource* resource) {
return static_cast<const Object*>(root)->template access_query_stream<Scene_Event_Stream_Tag>(
[&](std::span<const Event_Pointer> events) {
Event_Report_Batch result{resource};
result.reserve(events.size());
for (const auto& event : events) result.push_back(event);
return result;
});
};
}
template <Event_Object Event_Object_Type, typename... Arguments>
std::shared_ptr<Event_Object_Type> Scene::make_event(Arguments&&... arguments) {
+47
View File
@@ -1,7 +1,17 @@
#include "double_buffer/model.hpp"
#include <gtest/gtest.h>
#include <thread>
namespace {
struct Object_Buffer_Tag {};
struct Stream_Buffer_Tag {};
struct Stream_Object : double_buffer::Def<
Stream_Object, double_buffer::Root,
double_buffer::Mpmc_Triple_Buffer<Stream_Buffer_Tag, int>> {
struct Prop : Prev_Prop {};
struct State : Prev_State { bool operator==(const State&) const = default; };
struct Private : Prev_Private {};
};
using Stream = double_buffer::Impl<Stream_Object>;
struct Test_Object : double_buffer::Def<Test_Object, double_buffer::Root, double_buffer::Tagged_Buffer<Object_Buffer_Tag, int>> {
struct Prop : Prev_Prop {
int first{};
@@ -92,6 +102,43 @@ TEST(object_buffer, prop_publishes_and_keeps_incremental_baseline) {
EXPECT_EQ(object->data_for_test().pending->first, 17);
EXPECT_EQ(object->data_for_test().pending->second, 19);
}
TEST(object_buffer, mpmc_triple_buffer_rotates_three_queue_roles) {
auto object = build_object<Stream>();
object->submit_stream<Stream_Buffer_Tag>(11);
object->submit_stream<Stream_Buffer_Tag>(13);
object->exchange_stream<Stream_Buffer_Tag>();
object->access_rendering_stream<Stream_Buffer_Tag>([](std::span<int> rendering) {
ASSERT_EQ(rendering.size(), 2);
EXPECT_EQ(rendering[0], 11);
EXPECT_EQ(rendering[1], 13);
});
EXPECT_EQ(object->access_query_stream<Stream_Buffer_Tag>(
[](std::span<const int> query) { return query.size(); }), 0);
object->exchange_stream<Stream_Buffer_Tag>();
object->access_query_stream<Stream_Buffer_Tag>([](std::span<const int> query) {
ASSERT_EQ(query.size(), 2);
EXPECT_EQ(query[0], 11);
EXPECT_EQ(query[1], 13);
});
}
TEST(object_buffer, mpmc_triple_buffer_accepts_multiple_producers) {
auto object = build_object<Stream>();
std::vector<std::thread> producers;
for (int producer = 0; producer < 4; ++producer)
producers.emplace_back([&, producer] {
for (int index = 0; index < 100; ++index)
object->submit_stream<Stream_Buffer_Tag>(producer * 100 + index);
});
for (auto& producer : producers) producer.join();
object->exchange_stream<Stream_Buffer_Tag>();
EXPECT_EQ(object->access_rendering_stream<Stream_Buffer_Tag>(
[](std::span<int> rendering) { return rendering.size(); }), 400);
EXPECT_EQ(object->access_query_stream<Stream_Buffer_Tag>(
[](std::span<const int> query) { return query.size(); }), 0);
object->exchange_stream<Stream_Buffer_Tag>();
EXPECT_EQ(object->access_query_stream<Stream_Buffer_Tag>(
[](std::span<const int> query) { return query.size(); }), 400);
}
TEST(state_tag, callback_publishes_only_requested_layer) {
auto object = build_object<Object>();
int calls = 0;