Files
Renderive/Kernel/tests/renderive/threading/Threading_Contract_Test.cpp
T
2026-08-12 00:06:44 +08:00

605 lines
25 KiB
C++

#include <gtest/gtest.h>
#include <algorithm>
#include <atomic>
#include <cstdint>
#include <cmath>
#include <limits>
#include <memory>
#include <mutex>
#include <thread>
#include <vector>
#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<bool>& start) {
while (!start.load(std::memory_order_acquire)) {
std::this_thread::yield();
}
}
struct Threading_Test_Observation {};
struct Threading_Test_Observer_Data {
std::atomic<int> active{};
std::atomic<int> maximum{};
std::atomic<int> count{};
};
struct Threading_Test_Observer {
static constexpr bool enabled = true;
std::shared_ptr<Threading_Test_Observer_Data> data{std::make_shared<Threading_Test_Observer_Data>()};
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<Threading_Test_Observer, Threading_Test_Time_Source>;
struct Threading_Revision_Observer_Data {
std::mutex mutex;
std::vector<std::uint64_t> revisions;
};
struct Threading_Revision_Observer {
static constexpr bool enabled = true;
std::shared_ptr<Threading_Revision_Observer_Data> data{std::make_shared<Threading_Revision_Observer_Data>()};
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<Threading_Revision_Observer, Threading_Test_Time_Source>;
struct Threading_Test_State_Base {};
struct Threading_Test_State_Value {
Threading_Test_Frame frame;
};
using Threading_Test_Double_State = Double_State_Strategy<Threading_Test_State_Base, Threading_Test_State_Value, Atomic_Wait_Mutex>;
using Threading_Test_Triple_State = Triple_State_Strategy<Threading_Test_State_Base, Threading_Test_State_Value, Atomic_Wait_Mutex>;
struct Threading_Test_Scene_Renderable : Renderable_Base {
Threading_Test_Scene_Renderable(Scene_Base& scene, std::atomic<int>& render_count, std::atomic<std::uint64_t>& 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<int>* render_count;
std::atomic<std::uint64_t>* maximum_sequence;
};
struct Threading_Test_Task_Graph_Renderable : Renderable_Base {
Threading_Test_Task_Graph_Renderable(Scene_Base& scene, std::atomic<int>& 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<int>* 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<bool> start{};
std::vector<std::thread> 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<Threading_Test_Frame> 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<std::atomic<unsigned>[]>(frame_count);
std::atomic<bool> start{};
std::atomic<int> producers_done{};
std::atomic<int> consumed{};
std::atomic<int> invalid{};
std::vector<std::thread> 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<std::uint64_t>(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<std::thread> 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<std::uint64_t>(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<Threading_Test_Frame> strategy;
constexpr int producer_count = 4;
constexpr int frames_per_producer = 500;
constexpr int frame_count = producer_count * frames_per_producer;
std::atomic<bool> start{};
std::atomic<int> producers_done{};
std::atomic<bool> refresher_done{};
std::atomic<int> invalid{};
std::vector<std::thread> 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<std::uint64_t>(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<Threading_Test_Frame> 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<bool> start{};
std::atomic<int> producers_done{};
std::atomic<bool> updater_done{};
std::atomic<bool> configuration_done{};
std::atomic<int> invalid{};
std::vector<std::thread> 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<std::uint64_t>(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<std::uint64_t>(index), static_cast<std::uint64_t>(index), static_cast<std::uint64_t>(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<int, std::mutex, Threading_Revision_Observer_State> 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<bool> start{};
std::vector<std::thread> 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<Threading_Test_Frame> data;
constexpr int updater_count = 4;
constexpr int updates_per_thread = 1000;
constexpr int update_count = updater_count * updates_per_thread;
std::atomic<bool> start{};
std::atomic<int> updates_done{};
std::atomic<int> 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<std::thread> 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<std::uint64_t>(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<Threading_Test_Frame> data;
constexpr int updater_count = 4;
constexpr int updates_per_thread = 500;
constexpr int update_count = updater_count * updates_per_thread;
std::atomic<bool> start{};
std::atomic<int> updates_done{};
std::atomic<bool> discard_done{};
std::atomic<int> 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<std::uint64_t>::max());
std::this_thread::yield();
}
data.discard_before_time_ns(std::numeric_limits<std::uint64_t>::max());
discard_done.store(true, std::memory_order_release);
});
std::vector<std::thread> 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<std::uint64_t>(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<bool> start{};
std::atomic<bool> writer_done{};
std::atomic<int> invalid{};
std::thread writer([&] {
wait_start(start);
for (int index = 1; index <= update_count; ++index) {
const auto value = static_cast<std::uint64_t>(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<bool> start{};
std::atomic<bool> writer_done{};
std::atomic<int> invalid{};
std::thread writer([&] {
wait_start(start);
for (int index = 1; index <= update_count; ++index) {
const auto value = static_cast<std::uint64_t>(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<int> executed{};
std::atomic<std::uint64_t> maximum_sequence{};
auto renderable = std::make_shared<Threading_Test_Scene_Renderable>(scene, executed, maximum_sequence);
scene.attach_renderable(renderable);
std::atomic<bool> start{};
std::vector<std::thread> 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<int> executed{};
auto renderable = std::make_shared<Threading_Test_Task_Graph_Renderable>(scene, executed);
scene.attach_renderable(renderable);
std::atomic<bool> start{};
std::atomic<bool> 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<int, std::mutex, Observer_State<Frame_Strategy_Real_Time_Data_Observer>>;
struct Renderable : Renderable_Base {
explicit Renderable(Scene_Base& scene) : Renderable_Base(scene) {}
};
using Attachment = Attach_Real_Time_Data<Renderable, Data>;
auto data = std::make_shared<Data>();
auto scene = std::make_unique<Scene2D_Context<>>();
auto renderable = std::make_shared<Attachment>(With_Real_Time_Data(data), *scene);
std::atomic<bool> start{};
std::atomic<bool> stop{};
std::atomic<int> 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<Threading_Test_Frame>;
using Manual = Manual_Refresh_Strategy<Threading_Test_Frame>;
using Low_Latency = Low_Latency_Strategy<Threading_Test_Frame>;
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);
}