This commit is contained in:
2026-08-24 10:45:26 +08:00
parent 4f21169ade
commit cb9822074b
31 changed files with 363 additions and 536 deletions
@@ -244,8 +244,7 @@ struct Root::Builder {
Buffer_Build buffer_storage; /* build() 前暂存的普通 Buffer 初值。 */
Dependency_Graph_Build dependency_graph_storage; /* build() 前独立编辑并验证的依赖图。 */
template <typename... Args>
explicit Builder(Args&&... args) : object(new Object(std::forward<Args>(args)...)),
dependency_graph_storage(object->memory_resource()) {}
explicit Builder(Args&&... args) : object(new Object(std::forward<Args>(args)...)) {}
Final_Builder& set_pmr(Pmr pmr) {
object->set_pmr_resource(pmr);
return static_cast<Final_Builder&>(*this);
@@ -1,6 +1,7 @@
#pragma once
#include <concepts>
#include <algorithm>
#include <concurrentqueue-1.0.5/concurrentqueue.h>
#include <cstddef>
#include <functional>
@@ -52,6 +53,31 @@ public:
clear(*collecting_);
}
/*
* Accumulating streams promote the last published version into the current
* rendering version before appending the newly collected batch. Query keeps
* the previous complete version and producers remain isolated in collecting.
*/
void accumulate(std::size_t capacity)
requires std::copy_constructible<Value> {
accumulate_with([capacity](std::vector<Value>& values) {
if (values.size() <= capacity) return;
values.erase(values.begin(), values.end() -
static_cast<std::ptrdiff_t>(capacity));
});
}
template <typename Predicate>
requires std::copy_constructible<Value> &&
std::predicate<Predicate&, const Value&>
void accumulate(Predicate&& retain) {
accumulate_with([&](std::vector<Value>& values) {
std::erase_if(values, [&](const Value& value) {
return !std::invoke(retain, value);
});
});
}
template <typename Callback>
requires std::invocable<Callback&, std::span<Value>>
decltype(auto) access_rendering(Callback&& callback) {
@@ -82,6 +108,23 @@ private:
if (!queue.enqueue(std::move(value))) throw std::bad_alloc{};
}
template <typename Mutation>
void accumulate_with(Mutation&& mutate) {
std::unique_lock lock(exchange_lock_);
std::vector<Value> published;
published.reserve(query_->size_approx());
Value value;
while (query_->try_dequeue(value))
published.push_back(std::move(value));
std::vector<Value> accumulated{published.begin(), published.end()};
restore(*query_, published);
accumulated.reserve(accumulated.size() + rendering_->size_approx());
while (rendering_->try_dequeue(value))
accumulated.push_back(std::move(value));
std::invoke(std::forward<Mutation>(mutate), accumulated);
restore(*rendering_, accumulated);
}
template <typename Callback>
static decltype(auto) access(Queue& queue, Callback&& callback) {
std::vector<Value> values;
+10
View File
@@ -261,6 +261,16 @@ public:
void exchange_stream() {
data().mpmc_triple_buffer_storage.template get<Tag>().advance();
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag>
void accumulate_stream(std::size_t capacity) {
data().mpmc_triple_buffer_storage.template get<Tag>().accumulate(capacity);
}
template <detail::Mpmc_Triple_Buffer_Tag_In<Mpmc_Triple_Buffers> Tag,
typename Predicate>
void accumulate_stream(Predicate&& retain) {
data().mpmc_triple_buffer_storage.template get<Tag>().accumulate(
std::forward<Predicate>(retain));
}
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(
+22
View File
@@ -139,6 +139,28 @@ TEST(object_buffer, mpmc_triple_buffer_accepts_multiple_producers) {
EXPECT_EQ(object->access_query_stream<Stream_Buffer_Tag>(
[](std::span<const int> query) { return query.size(); }), 400);
}
TEST(object_buffer, mpmc_triple_buffer_accumulates_published_history_with_capacity) {
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->accumulate_stream<Stream_Buffer_Tag>(3);
object->submit_stream<Stream_Buffer_Tag>(17);
object->submit_stream<Stream_Buffer_Tag>(19);
object->exchange_stream<Stream_Buffer_Tag>();
object->accumulate_stream<Stream_Buffer_Tag>(3);
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);
});
object->access_rendering_stream<Stream_Buffer_Tag>([](std::span<int> rendering) {
ASSERT_EQ(rendering.size(), 3);
EXPECT_EQ(rendering[0], 13);
EXPECT_EQ(rendering[1], 17);
EXPECT_EQ(rendering[2], 19);
});
}
TEST(state_tag, callback_publishes_only_requested_layer) {
auto object = build_object<Object>();
int calls = 0;