仍然闪烁

This commit is contained in:
2026-08-02 19:26:50 +08:00
parent 9a39f88427
commit 36d63efc7b
10 changed files with 650 additions and 128 deletions
+22 -2
View File
@@ -1,6 +1,7 @@
#include <algorithm>
#include <utility>
#include "../base/Memory.h"
#include "../render/Frame_Prepare_Subflow.h"
#include "Plot_Render_Context.h"
#include "../plot/Plot_Core.h"
#include "../plot/Plot_Core_Access.h"
@@ -407,9 +408,28 @@ void Renderable::prepare_data(const Render_Frame_Snapshot& snapshot, const Rende
if (!frame_view)
return;
d_ptr->prepare_data(snapshot, frame_view);
if (snapshot.frame_update_states)
d_ptr->drain_committed_update_states(*snapshot.frame_update_states);
}
void build_renderable_prepare(
Renderable& renderable,
Render_Task_Subflow& subflow,
const Render_Frame_Snapshot& snapshot,
const Renderable_Frame_View& frame_view) {
if (renderable.get_retiring() || !frame_view || !renderable.d_ptr)
return;
if (auto* source = dynamic_cast<Frame_Prepare_Subtask_Source*>(renderable.d_ptr)) {
if (source->build_prepare_subtasks(subflow, snapshot, frame_view))
return;
}
renderable.prepare_data(snapshot, frame_view);
}
void drain_renderable_prepare_updates(Renderable& renderable, std::pmr::vector<std::shared_ptr<Update_State>>& output) {
if (renderable.get_retiring() || !renderable.d_ptr)
return;
renderable.d_ptr->drain_committed_update_states(output);
}
bool Renderable::select_test(const PointF& pos) {
if (get_retiring())
return false;
+3
View File
@@ -24,6 +24,7 @@
namespace renderive {
class Plot_Core;
class Renderable;
class Render_Task_Subflow;
struct Plot_Render_Context;
template <typename Owner, Render_Data_Type Data, Render_State_Type State, Input_Data_Type Input>
struct Renderable_Binding;
@@ -88,6 +89,8 @@ public:
[[nodiscard]] bool get_retiring() const;
virtual ~Renderable();
private:
friend void build_renderable_prepare(Renderable& renderable, Render_Task_Subflow& subflow, const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view);
friend void drain_renderable_prepare_updates(Renderable& renderable, std::pmr::vector<std::shared_ptr<Update_State>>& output);
template <typename Owner, Render_Data_Type Data, Render_State_Type State, Input_Data_Type Input>
friend struct Renderable_Binding;
template <typename That>
+151 -36
View File
@@ -9,37 +9,42 @@
#include <thread>
namespace renderive {
struct Render_Executor_Metrics;
struct Frame_Task_Counters {
std::atomic_uint32_t task_count{0};
std::atomic_uint32_t active_count{0};
std::atomic_uint32_t peak_parallelism{0};
std::atomic_uint64_t task_duration_total_ns{0};
};
struct Taskflow_Graph_Impl {
explicit Taskflow_Graph_Impl(tf::Taskflow& taskflow) : taskflow(taskflow) {}
Taskflow_Graph_Impl(
tf::Taskflow& taskflow,
std::shared_ptr<Render_Executor_Metrics> metrics,
std::shared_ptr<Frame_Task_Counters> frame_counters)
: taskflow(taskflow),
metrics(std::move(metrics)),
frame_counters(std::move(frame_counters)) {}
tf::Taskflow& taskflow;
std::shared_ptr<Render_Executor_Metrics> metrics;
std::shared_ptr<Frame_Task_Counters> frame_counters;
std::vector<tf::Task> tasks;
};
Render_Task_Graph::Render_Task_Graph(void* impl) noexcept : impl(impl) {}
Render_Task_Graph_Node Render_Task_Graph::emplace(Task work) {
auto* graph = static_cast<Taskflow_Graph_Impl*>(impl);
if (!graph || !work)
return {};
auto work_ptr = std::make_shared<Task>(std::move(work));
auto task = graph->taskflow.emplace([work_ptr]() mutable {
if (work_ptr && *work_ptr)
(*work_ptr)();
});
graph->tasks.push_back(task);
return {static_cast<std::uint32_t>(graph->tasks.size() - 1)};
}
void Render_Task_Graph::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) {
auto* graph = static_cast<Taskflow_Graph_Impl*>(impl);
if (!graph || !before || !after)
return;
if (before.index >= graph->tasks.size() || after.index >= graph->tasks.size())
return;
graph->tasks[before.index].precede(graph->tasks[after.index]);
}
namespace {
struct Taskflow_Subflow_Impl {
Taskflow_Subflow_Impl(
tf::Subflow& subflow,
std::shared_ptr<Render_Executor_Metrics> metrics,
std::shared_ptr<Frame_Task_Counters> frame_counters)
: subflow(subflow),
metrics(std::move(metrics)),
frame_counters(std::move(frame_counters)) {}
tf::Subflow& subflow;
std::shared_ptr<Render_Executor_Metrics> metrics;
std::shared_ptr<Frame_Task_Counters> frame_counters;
std::vector<tf::Task> tasks;
};
struct Render_Executor_Metrics {
std::size_t worker_count{};
@@ -89,7 +94,113 @@ void add_kind_count(Render_Executor_Metrics& metrics, Render_Task_Kind kind) {
}
}
} // namespace
void begin_business_task(Render_Executor_Metrics* metrics, Frame_Task_Counters* counters, Render_Task_Kind kind) {
if (metrics)
add_kind_count(*metrics, kind);
if (!counters)
return;
counters->task_count.fetch_add(1, std::memory_order_relaxed);
std::uint32_t active = counters->active_count.fetch_add(1, std::memory_order_acq_rel) + 1;
std::uint32_t peak = counters->peak_parallelism.load(std::memory_order_acquire);
while (peak < active && !counters->peak_parallelism.compare_exchange_weak(peak, active, std::memory_order_acq_rel, std::memory_order_acquire)) {}
}
void end_business_task(Frame_Task_Counters* counters, std::uint64_t begin_ns) {
if (!counters)
return;
std::uint64_t end_ns = steady_now_ns();
counters->task_duration_total_ns.fetch_add(end_ns > begin_ns ? end_ns - begin_ns : 0, std::memory_order_relaxed);
counters->active_count.fetch_sub(1, std::memory_order_acq_rel);
}
Task wrap_business_task(
Render_Task_Kind kind,
Task work,
std::shared_ptr<Render_Executor_Metrics> metrics,
std::shared_ptr<Frame_Task_Counters> counters) {
return Task([kind, work = std::move(work), metrics = std::move(metrics), counters = std::move(counters)]() mutable {
std::uint64_t begin_ns = steady_now_ns();
begin_business_task(metrics.get(), counters.get(), kind);
if (work)
work();
end_business_task(counters.get(), begin_ns);
});
}
Render_Task_Subflow::Render_Task_Subflow(void* impl) noexcept : impl(impl) {}
Render_Task_Graph_Node Render_Task_Subflow::emplace(Render_Task_Kind kind, Task work) {
auto* graph = static_cast<Taskflow_Subflow_Impl*>(impl);
if (!graph || !work)
return {};
auto wrapped = std::make_shared<Task>(wrap_business_task(kind, std::move(work), graph->metrics, graph->frame_counters));
auto task = graph->subflow.emplace([wrapped]() mutable {
if (wrapped && *wrapped)
(*wrapped)();
});
graph->tasks.push_back(task);
return {static_cast<std::uint32_t>(graph->tasks.size() - 1)};
}
void Render_Task_Subflow::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) {
auto* graph = static_cast<Taskflow_Subflow_Impl*>(impl);
if (!graph || !before || !after)
return;
if (before.index >= graph->tasks.size() || after.index >= graph->tasks.size())
return;
graph->tasks[before.index].precede(graph->tasks[after.index]);
}
void Render_Task_Subflow::join() {
if (joined)
return;
auto* graph = static_cast<Taskflow_Subflow_Impl*>(impl);
if (!graph)
return;
graph->subflow.join();
joined = true;
}
Render_Task_Graph::Render_Task_Graph(void* impl) noexcept : impl(impl) {}
Render_Task_Graph_Node Render_Task_Graph::emplace(Render_Task_Kind kind, Task work) {
auto* graph = static_cast<Taskflow_Graph_Impl*>(impl);
if (!graph || !work)
return {};
auto wrapped = std::make_shared<Task>(wrap_business_task(kind, std::move(work), graph->metrics, graph->frame_counters));
auto task = graph->taskflow.emplace([wrapped]() mutable {
if (wrapped && *wrapped)
(*wrapped)();
});
graph->tasks.push_back(task);
return {static_cast<std::uint32_t>(graph->tasks.size() - 1)};
}
Render_Task_Graph_Node Render_Task_Graph::emplace_subflow(Render_Task_Kind kind, Subflow_Builder builder) {
auto* graph = static_cast<Taskflow_Graph_Impl*>(impl);
if (!graph || !builder)
return {};
auto task = graph->taskflow.emplace([kind, builder = std::move(builder), metrics = graph->metrics, counters = graph->frame_counters](tf::Subflow& raw_subflow) mutable {
std::uint64_t begin_ns = steady_now_ns();
begin_business_task(metrics.get(), counters.get(), kind);
Taskflow_Subflow_Impl subflow_impl(raw_subflow, metrics, counters);
Render_Task_Subflow subflow(&subflow_impl);
builder(subflow);
subflow.join();
end_business_task(counters.get(), begin_ns);
});
graph->tasks.push_back(task);
return {static_cast<std::uint32_t>(graph->tasks.size() - 1)};
}
void Render_Task_Graph::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) {
auto* graph = static_cast<Taskflow_Graph_Impl*>(impl);
if (!graph || !before || !after)
return;
if (before.index >= graph->tasks.size() || after.index >= graph->tasks.size())
return;
graph->tasks[before.index].precede(graph->tasks[after.index]);
}
class Render_Executor_Observer final : public tf::ObserverInterface {
public:
@@ -209,7 +320,6 @@ public:
private:
void record_submission(Render_Executor_Task& task) {
metrics->submitted_tasks.fetch_add(1, std::memory_order_relaxed);
add_kind_count(*metrics, task.kind);
if (task.frame_job)
metrics->queued_frame_jobs.fetch_add(1, std::memory_order_relaxed);
if (task.frame_stat)
@@ -220,6 +330,7 @@ private:
auto task_ptr = std::make_shared<Render_Executor_Task>(std::move(task));
auto completion = std::make_shared<Task>(std::move(task_ptr->completion));
auto run_record = std::make_shared<Task_Run_Record>();
auto frame_counters = std::make_shared<Frame_Task_Counters>();
auto topology = std::make_shared<tf::Taskflow>();
auto entry = topology->emplace([this, enqueue_ns, task_ptr, run_record]() mutable {
Render_Executor_Task& task = *task_ptr;
@@ -234,7 +345,7 @@ private:
metrics->active_frame_jobs.fetch_add(1, std::memory_order_relaxed);
}
});
auto exit = topology->emplace([this, task_ptr, run_record]() mutable {
auto exit = topology->emplace([this, task_ptr, run_record, frame_counters]() mutable {
Render_Executor_Task& task = *task_ptr;
std::uint64_t end_ns = steady_now_ns();
std::uint64_t begin_ns = run_record->begin_ns;
@@ -242,11 +353,14 @@ private:
std::uint64_t run_ns = end_ns > begin_ns ? end_ns - begin_ns : 0;
if (task.frame_stat) {
Frame_Worker_Stats& stat = *task.frame_stat;
stat.task_count += std::max<std::uint32_t>(1, task.logical_task_count);
stat.peak_parallelism = static_cast<std::uint32_t>(std::max<std::uint64_t>(stat.peak_parallelism, metrics->peak_concurrency.load(std::memory_order_acquire)));
std::uint32_t task_count = frame_counters ? frame_counters->task_count.load(std::memory_order_acquire) : 0;
std::uint32_t peak_parallelism = frame_counters ? frame_counters->peak_parallelism.load(std::memory_order_acquire) : 0;
std::uint64_t task_duration_total_ns = frame_counters ? frame_counters->task_duration_total_ns.load(std::memory_order_acquire) : 0;
stat.task_count += task_count;
stat.peak_parallelism = std::max(stat.peak_parallelism, peak_parallelism);
stat.queue_wait_total_ns += queue_wait_ns;
stat.queue_wait_max_ns = std::max(stat.queue_wait_max_ns, queue_wait_ns);
stat.worker_run_total_ns += run_ns;
stat.worker_run_total_ns += task_duration_total_ns;
stat.parallel_stage_wall_ns += run_ns;
stat.executor_at_finish = snapshot();
}
@@ -254,16 +368,17 @@ private:
metrics->active_frame_jobs.fetch_sub(1, std::memory_order_relaxed);
});
if (task_ptr->build_graph) {
Taskflow_Graph_Impl graph_impl(*topology);
Taskflow_Graph_Impl graph_impl(*topology, metrics, frame_counters);
graph_impl.tasks.push_back(entry);
graph_impl.tasks.push_back(exit);
Render_Task_Graph graph(&graph_impl);
task_ptr->build_graph(graph, {0}, {1});
}
else {
auto work = topology->emplace([task_ptr]() mutable {
if (task_ptr->work)
task_ptr->work();
auto wrapped = std::make_shared<Task>(wrap_business_task(task_ptr->kind, std::move(task_ptr->work), metrics, frame_counters));
auto work = topology->emplace([wrapped]() mutable {
if (wrapped && *wrapped)
(*wrapped)();
});
entry.precede(work);
work.precede(exit);
+14 -2
View File
@@ -12,10 +12,23 @@ struct Render_Task_Graph_Node {
return index != UINT32_MAX;
}
};
class LIB_DECL Render_Task_Subflow {
public:
explicit Render_Task_Subflow(void* impl) noexcept;
[[nodiscard]] Render_Task_Graph_Node emplace(Render_Task_Kind kind, Task work);
void precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after);
void join();
private:
void* impl{};
bool joined{};
};
class LIB_DECL Render_Task_Graph {
public:
using Subflow_Builder = std::function<void(Render_Task_Subflow&)>;
explicit Render_Task_Graph(void* impl) noexcept;
[[nodiscard]] Render_Task_Graph_Node emplace(Task work);
[[nodiscard]] Render_Task_Graph_Node emplace(Render_Task_Kind kind, Task work);
[[nodiscard]] Render_Task_Graph_Node emplace_subflow(Render_Task_Kind kind, Subflow_Builder builder);
void precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after);
private:
void* impl{};
@@ -29,7 +42,6 @@ struct Render_Executor_Task {
Graph_Builder build_graph;
Task completion;
std::shared_ptr<Frame_Worker_Stats> frame_stat;
std::uint32_t logical_task_count{1};
bool frame_job{};
};
class Render_Executor_Private;
+40 -11
View File
@@ -6,6 +6,7 @@
#include "../architecture/Update_Completion.h"
#include "../base/global.h"
#include "../render/Canvas.h"
#include "../render/Frame_Prepare_Subflow.h"
#include "../render/Render_Frame_Snapshot.h"
#include <algorithm>
#include <memory>
@@ -21,6 +22,7 @@ struct Frame_Render_Job {
Render_Frame_Snapshot render_snapshot;
Renderable_Frame_Tree renderables;
std::shared_ptr<Frame_Worker_Stats> worker_stat;
std::atomic_uint64_t prepared_output_count{0};
};
class Plot_Core_Private final : public std::enable_shared_from_this<Plot_Core_Private>, public Frame_Flow_Host {
public:
@@ -92,17 +94,25 @@ public:
Render_Task_Kind::Frame,
Task{},
[job](Render_Task_Graph& graph, Render_Task_Graph_Node entry, Render_Task_Graph_Node exit) {
auto prepare = graph.emplace(Task([job]() {
job->owner->prepare_frame_job(job);
auto prepare_begin = graph.emplace(Render_Task_Kind::Frame, Task([job]() {
job->owner->begin_prepare_frame_job(job);
}));
auto draw = graph.emplace(Task([job]() {
auto prepare_dispatch = graph.emplace_subflow(Render_Task_Kind::Frame, [job](Render_Task_Subflow& subflow) {
job->owner->dispatch_prepare_subflow(job, subflow);
});
auto prepare_end = graph.emplace(Render_Task_Kind::Frame, Task([job]() {
job->owner->end_prepare_frame_job(job);
}));
auto draw = graph.emplace(Render_Task_Kind::Primitive, Task([job]() {
job->owner->draw_frame_job(job);
}));
auto finalize = graph.emplace(Task([job]() {
auto finalize = graph.emplace(Render_Task_Kind::Frame, Task([job]() {
job->owner->finalize_frame_job(job);
}));
graph.precede(entry, prepare);
graph.precede(prepare, draw);
graph.precede(entry, prepare_begin);
graph.precede(prepare_begin, prepare_dispatch);
graph.precede(prepare_dispatch, prepare_end);
graph.precede(prepare_end, draw);
graph.precede(draw, finalize);
graph.precede(finalize, exit);
},
@@ -112,7 +122,6 @@ public:
});
}),
job->worker_stat,
3,
true
});
if (submitted)
@@ -128,18 +137,35 @@ public:
static bool frame_job_valid(const std::shared_ptr<Frame_Render_Job>& job) {
return job && job->buffer && job->lifecycle && !job->render_snapshot.destroying();
}
void prepare_frame_job(const std::shared_ptr<Frame_Render_Job>& job) {
void begin_prepare_frame_job(const std::shared_ptr<Frame_Render_Job>& job) {
if (!frame_job_valid(job))
return;
Frame_Lifecycle_Record& frame = *job->lifecycle;
frame.prepare_begin_ns = steady_now_ns();
job->prepared_output_count.store(0, std::memory_order_release);
}
void dispatch_prepare_subflow(const std::shared_ptr<Frame_Render_Job>& job, Render_Task_Subflow& subflow) {
if (!frame_job_valid(job) || !job->lifecycle->prepare_begin_ns)
return;
for (const auto& node : job->renderables) {
const auto& renderable = node.renderable;
if (!renderable || !renderable->should_prepare())
continue;
renderable->prepare_data(job->render_snapshot, node.frame_view);
++frame.prepared_output_count;
build_renderable_prepare(*renderable, subflow, job->render_snapshot, node.frame_view);
job->prepared_output_count.fetch_add(1, std::memory_order_relaxed);
}
}
void end_prepare_frame_job(const std::shared_ptr<Frame_Render_Job>& job) {
if (!frame_job_valid(job) || !job->lifecycle->prepare_begin_ns)
return;
Frame_Lifecycle_Record& frame = *job->lifecycle;
if (job->render_snapshot.frame_update_states) {
for (const auto& node : job->renderables) {
if (node.renderable)
drain_renderable_prepare_updates(*node.renderable, *job->render_snapshot.frame_update_states);
}
}
frame.prepared_output_count = job->prepared_output_count.load(std::memory_order_acquire);
frame.prepare_end_ns = steady_now_ns();
if (job->owner && job->owner->context)
job->owner->context->first_prepare_data.store(false, std::memory_order_release);
@@ -169,7 +195,10 @@ public:
void execute_frame_job(const std::shared_ptr<Frame_Render_Job>& job) {
if (!job || !job->buffer || !job->lifecycle || job->render_snapshot.destroying())
return;
prepare_frame_job(job);
begin_prepare_frame_job(job);
Render_Task_Subflow no_subflow(nullptr);
dispatch_prepare_subflow(job, no_subflow);
end_prepare_frame_job(job);
draw_frame_job(job);
finalize_frame_job(job);
}
+67 -23
View File
@@ -19,6 +19,7 @@
#include "../base/Memory.h"
#include "../base/global.h"
#include "../render/Canvas.h"
#include "../render/Frame_Prepare_Subflow.h"
#include "../render/Image.h"
#include "../render/Render_Frame_Snapshot.h"
#include "Psc_Cpp_Core/Base/RingBuffer.hpp"
@@ -143,15 +144,25 @@ struct Waterfall_Image_Tile {
Image image;
};
struct Waterfall_Image_Job {
explicit Waterfall_Image_Job(std::shared_ptr<Waterfall_Image_Request> request, int columns, int rows)
: request(std::move(request)), tile_columns(columns), tile_rows(rows), remaining(columns * rows) {
explicit Waterfall_Image_Job(std::shared_ptr<const Waterfall_Image_Request> request, int columns, int rows)
: request(std::move(request)),
update_states(memory_resource(Memory_Domain::Update_Completion)),
tile_columns(columns),
tile_rows(rows) {
tiles.resize(static_cast<std::size_t>(columns * rows));
}
std::shared_ptr<Waterfall_Image_Request> request;
~Waterfall_Image_Job() {
for (auto& state : update_states) {
if (state)
state->complete_all(std::make_error_code(std::errc::operation_canceled), Update_Outcome::Cancelled);
}
update_states.clear();
}
std::shared_ptr<const Waterfall_Image_Request> request;
std::pmr::vector<std::shared_ptr<Update_State>> update_states;
int tile_columns{};
int tile_rows{};
std::vector<Waterfall_Image_Tile> tiles;
std::atomic<int> remaining;
};
static void complete_waterfall_request_states(Waterfall_Image_Request& request, std::error_code error, Update_Outcome outcome) {
for (auto& state : request.update_states) {
@@ -186,6 +197,19 @@ static int waterfall_tile_row_count(int height) {
static constexpr int Tile_Height = 64;
return std::max(1, (height + Tile_Height - 1) / Tile_Height);
}
static std::shared_ptr<Waterfall_Image_Job> create_waterfall_image_job(std::shared_ptr<Waterfall_Image_Request> request) {
if (!request || request->key.width <= 0 || request->key.height <= 0)
return {};
int columns = waterfall_tile_column_count(request->key.width);
int rows = waterfall_tile_row_count(request->key.height);
std::pmr::vector<std::shared_ptr<Update_State>> update_states = std::move(request->update_states);
auto job = std::make_shared<Waterfall_Image_Job>(
std::shared_ptr<const Waterfall_Image_Request>(std::move(request)),
columns,
rows);
job->update_states = std::move(update_states);
return job;
}
static void write_waterfall_tile(Waterfall_Image_Job& job, int tile_index) {
const Waterfall_Image_Request& request = *job.request;
int tile_x = tile_index % job.tile_columns;
@@ -222,25 +246,15 @@ static void write_waterfall_tile(Waterfall_Image_Job& job, int tile_index) {
}
}
}
static void write_waterfall_tiles(Waterfall_Image_Job& job) {
int tile_count = static_cast<int>(job.tiles.size());
for (int tile_index = 0; tile_index < tile_count; ++tile_index)
write_waterfall_tile(job, tile_index);
}
static Waterfall_Image_Snapshot build_waterfall_image_snapshot(const std::shared_ptr<Waterfall_Image_Request>& request_ptr) {
Waterfall_Image_Request& request = *request_ptr;
Waterfall_Image_Job job(
request_ptr,
waterfall_tile_column_count(request.key.width),
waterfall_tile_row_count(request.key.height));
write_waterfall_tiles(job);
static Waterfall_Image_Snapshot compose_waterfall_image(Waterfall_Image_Job& job) {
const Waterfall_Image_Request& request = *job.request;
Waterfall_Image_Snapshot snapshot;
snapshot.key = request.key;
snapshot.frequency_range = request.draw_frequency_range;
snapshot.time_range = request.draw_time_range;
snapshot.source_rect = request.draw_source_rect;
snapshot.image.resize(request.key.width, request.key.height);
for (const auto& tile : job.tiles) {
for (const Waterfall_Image_Tile& tile : job.tiles) {
for (int y = 0; y < tile.image.height(); ++y) {
Pixel* target = snapshot.image.row(tile.y + y);
const Pixel* source = tile.image.row(y);
@@ -248,10 +262,10 @@ static Waterfall_Image_Snapshot build_waterfall_image_snapshot(const std::shared
std::memcpy(target + tile.x, source, static_cast<std::size_t>(tile.image.width()) * sizeof(Pixel));
}
}
snapshot.update_states = std::move(request.update_states);
snapshot.update_states = std::move(job.update_states);
return snapshot;
}
struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State, Waterfall_Input_Data>, Hit_Testable, Hover_Interactive {
struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State, Waterfall_Input_Data>, Hit_Testable, Hover_Interactive, Frame_Prepare_Subtask_Source {
static constexpr std::size_t Row_Queue_Capacity = 256;
Waterfall_Ring_Buffer ring_buffer;
std::weak_ptr<Color_Bar> color_bar;
@@ -357,6 +371,13 @@ struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State,
return true;
}
void prepare_data(const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) override {
prepare_waterfall(snapshot, frame_view, nullptr);
}
bool build_prepare_subtasks(Render_Task_Subflow& subflow, const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) override {
prepare_waterfall(snapshot, frame_view, &subflow);
return true;
}
void prepare_waterfall(const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view, Render_Task_Subflow* subflow) {
Waterfall_Render_State* s = render_state(frame_view);
if (!s)
return;
@@ -424,9 +445,9 @@ struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State,
break;
}
}
update_image_item(snapshot, *timeline, s);
update_image_item(snapshot, *timeline, s, subflow);
}
void update_image_item(const Render_Frame_Snapshot& snapshot, const Timeline_Stream_Snapshot& timeline, Waterfall_Render_State* s) {
void update_image_item(const Render_Frame_Snapshot& snapshot, const Timeline_Stream_Snapshot& timeline, Waterfall_Render_State* s, Render_Task_Subflow* subflow) {
if (!s)
return;
auto frequency_axis = s->frequency_axis.lock();
@@ -469,7 +490,8 @@ struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State,
};
bool should_submit = !submitted_image_key || !(*submitted_image_key == key) || !pending_image_update_states.empty();
if (should_submit) {
submitted_image_key = key;
if (!subflow)
return;
auto request = std::make_shared<Waterfall_Image_Request>();
request->key = key;
request->color_map = color_scale->color_map;
@@ -497,7 +519,29 @@ struct Waterfall_Private : Typed_Render_Data<Waterfall, Waterfall_Render_State,
}
request->update_states = std::move(pending_image_update_states);
pending_image_update_states = std::pmr::vector<std::shared_ptr<Update_State>>(memory_resource(Memory_Domain::Update_Completion));
image_snapshot = build_waterfall_image_snapshot(request);
auto job = create_waterfall_image_job(std::move(request));
if (!job || job->tiles.empty())
return;
std::vector<Render_Task_Graph_Node> tile_tasks;
tile_tasks.reserve(job->tiles.size());
for (std::size_t i = 0; i < job->tiles.size(); ++i) {
tile_tasks.push_back(subflow->emplace(Render_Task_Kind::Waterfall, Task([job, i]() {
write_waterfall_tile(*job, static_cast<int>(i));
})));
}
auto compose = subflow->emplace(Render_Task_Kind::Compose, Task([this, job]() {
if (q() && q()->get_retiring())
return;
Waterfall_Image_Snapshot result = compose_waterfall_image(*job);
if (q() && q()->get_retiring()) {
complete_waterfall_result_states(result, std::make_error_code(std::errc::operation_canceled), Update_Outcome::Cancelled);
return;
}
image_snapshot = std::move(result);
submitted_image_key = image_snapshot->key;
}));
for (auto tile : tile_tasks)
subflow->precede(tile, compose);
}
}
void draw_image(Canvas& canvas, const Render_Frame_Snapshot& snapshot, Waterfall_Render_State* s, Abs_Axis* h_axis, Time_Axis* v_axis) {
+224 -40
View File
@@ -11,7 +11,9 @@
#include <system_error>
#include <variant>
#include "../base/Memory.h"
#include "../base/global.h"
#include "../render/Canvas.h"
#include "../render/Frame_Prepare_Subflow.h"
#include "../render/Render_Frame_Snapshot.h"
namespace renderive {
enum class Curve_Input_Mode : std::uint8_t {
@@ -137,16 +139,46 @@ struct Curve_Geometry_Request {
std::shared_ptr<const Curve_Source> source;
std::pmr::vector<std::shared_ptr<Update_State>> update_states;
};
struct Curve_Geometry_Partition {
std::size_t source_begin{};
std::size_t source_end{};
int pixel_begin{};
int pixel_end{};
bool joins_previous{};
};
struct Curve_Geometry_Partial_Result {
Curve_Geometry_Partial_Result()
: segments(memory_resource(Memory_Domain::Curve)) {}
std::pmr::vector<Curve_Geometry_Result::Segment> segments;
std::size_t source_begin{};
std::size_t source_end{};
bool begins_with_continuous_point{};
bool ends_with_continuous_point{};
bool joins_previous{};
};
struct Curve_Geometry_Job {
explicit Curve_Geometry_Job(std::shared_ptr<const Curve_Geometry_Request> request, Curve_Geometry_Key key)
: request(std::move(request)),
key(std::move(key)),
update_states(memory_resource(Memory_Domain::Update_Completion)) {}
~Curve_Geometry_Job() {
for (auto& state : update_states) {
if (state)
state->complete_all(std::make_error_code(std::errc::operation_canceled), Update_Outcome::Cancelled);
}
update_states.clear();
}
std::shared_ptr<const Curve_Geometry_Request> request;
Curve_Geometry_Key key;
std::vector<Curve_Geometry_Partition> partitions;
std::vector<Curve_Geometry_Partial_Result> results;
std::pmr::vector<std::shared_ptr<Update_State>> update_states;
};
static constexpr std::size_t Min_Curve_Items_Per_Task = 4096;
static constexpr std::size_t Max_Curve_Partition_Count = 8;
static Curve_Axis_Geometry_Key curve_axis_key(const Axis_Frame_Snapshot& axis) {
return {axis.orientation, axis.coord_range, axis.pixel_range};
}
static void complete_curve_update_states(Curve_Geometry_Request& request, std::error_code error, Update_Outcome outcome) {
for (auto& state : request.update_states) {
if (state)
state->complete_all(error, outcome);
}
request.update_states.clear();
}
static void complete_curve_result_states(Curve_Geometry_Result& result, std::error_code error, Update_Outcome outcome) {
for (const auto& state : result.update_states) {
if (state)
@@ -213,37 +245,166 @@ static void append_curve_sample_run(Curve_Geometry_Result& result, const Curve_G
result.segments.back().points = std::move(points);
}
static void build_curve_geometry(Curve_Geometry_Result& result, const Curve_Geometry_Request& request) {
result.segments.clear();
std::size_t resolve_curve_partition_count(std::size_t source_count, std::size_t worker_count) {
if (source_count == 0)
return 0;
if (worker_count <= 1 || source_count < Min_Curve_Items_Per_Task)
return 1;
std::size_t by_size = (source_count + Min_Curve_Items_Per_Task - 1) / Min_Curve_Items_Per_Task;
return std::max<std::size_t>(1, std::min({worker_count, by_size, Max_Curve_Partition_Count}));
}
static std::size_t curve_request_source_count(const Curve_Geometry_Request& request) {
if (!request.source)
return;
return 0;
if (const auto* data_points = std::get_if<Curve_Point_Source>(request.source.get()))
return data_points->values.size();
if (const auto* samples = std::get_if<Curve_Sample_Source>(request.source.get()))
return samples->values.size();
return 0;
}
static bool curve_sample_valid(const Curve_Geometry_Request& request, std::size_t index) {
if (!request.source)
return false;
if (const auto* data_points = std::get_if<Curve_Point_Source>(request.source.get())) {
result.segments.reserve(data_points->values.size());
Curve_Geometry_Result::Segment* segment = nullptr;
for (const PointF& point : data_points->values)
append_curve_point_segment(result, segment, request.mapping, point);
if (index >= data_points->values.size())
return false;
PointF point = data_points->values[index];
return finite_point(point) && finite_point(request.mapping.map(point));
}
if (const auto* samples = std::get_if<Curve_Sample_Source>(request.source.get())) {
if (index >= samples->values.size())
return false;
double value = samples->values[index];
if (!std::isfinite(value))
return false;
double coord = source_index_to_coordinate(request.data_range, static_cast<int>(samples->values.size()), static_cast<double>(index));
return finite_point(request.mapping.map(coord, value));
}
return false;
}
static std::shared_ptr<Curve_Geometry_Job> create_curve_geometry_job(std::shared_ptr<Curve_Geometry_Request> request, Curve_Geometry_Key key, std::size_t worker_count) {
if (!request)
return {};
std::size_t source_count = curve_request_source_count(*request);
std::size_t partition_count = resolve_curve_partition_count(source_count, worker_count);
std::pmr::vector<std::shared_ptr<Update_State>> update_states = std::move(request->update_states);
auto job = std::make_shared<Curve_Geometry_Job>(
std::shared_ptr<const Curve_Geometry_Request>(std::move(request)),
std::move(key));
job->update_states = std::move(update_states);
job->partitions.reserve(partition_count);
job->results.resize(partition_count);
for (std::size_t i = 0; i < partition_count; ++i) {
std::size_t begin = source_count * i / partition_count;
std::size_t end = source_count * (i + 1) / partition_count;
Curve_Geometry_Partition partition;
partition.source_begin = begin;
partition.source_end = end;
partition.joins_previous = begin > 0 && curve_sample_valid(*job->request, begin - 1) && curve_sample_valid(*job->request, begin);
job->partitions.push_back(partition);
}
return job;
}
static void build_curve_partition(Curve_Geometry_Job& job, std::size_t partition_index) {
if (!job.request || partition_index >= job.partitions.size() || partition_index >= job.results.size())
return;
const Curve_Geometry_Request& request = *job.request;
const Curve_Geometry_Partition& partition = job.partitions[partition_index];
Curve_Geometry_Partial_Result partial;
partial.source_begin = partition.source_begin;
partial.source_end = partition.source_end;
partial.joins_previous = partition.joins_previous;
Curve_Geometry_Result result;
result.mapping = request.mapping;
if (!request.source || partition.source_begin >= partition.source_end) {
job.results[partition_index] = std::move(partial);
return;
}
const auto* samples = std::get_if<Curve_Sample_Source>(request.source.get());
if (!samples || samples->values.size() < 2 || request.data_range.length() == 0.0)
if (const auto* data_points = std::get_if<Curve_Point_Source>(request.source.get())) {
result.segments.reserve(partition.source_end - partition.source_begin);
Curve_Geometry_Result::Segment* segment = nullptr;
for (std::size_t i = partition.source_begin; i < partition.source_end && i < data_points->values.size(); ++i)
append_curve_point_segment(result, segment, request.mapping, data_points->values[i]);
}
else if (const auto* samples = std::get_if<Curve_Sample_Source>(request.source.get())) {
if (samples->values.size() >= 2 && request.data_range.length() != 0.0) {
std::size_t compute_begin = partition.source_begin;
if (partition.joins_previous && compute_begin > 0)
--compute_begin;
int run_start = -1;
std::size_t end = std::min(partition.source_end, samples->values.size());
for (std::size_t i = compute_begin; i < end; ++i) {
double value = samples->values[i];
if (std::isfinite(value)) {
if (run_start < 0)
run_start = static_cast<int>(i);
continue;
}
if (run_start >= 0)
append_curve_sample_run(result, request, *samples, run_start, static_cast<int>(i) - 1);
run_start = -1;
}
if (run_start >= 0)
append_curve_sample_run(result, request, *samples, run_start, static_cast<int>(end) - 1);
}
}
partial.begins_with_continuous_point = partition.source_begin < partition.source_end && curve_sample_valid(request, partition.source_begin);
partial.ends_with_continuous_point = partition.source_end > partition.source_begin && curve_sample_valid(request, partition.source_end - 1);
partial.segments = std::move(result.segments);
job.results[partition_index] = std::move(partial);
}
static void append_curve_segment(Curve_Geometry_Result& result, Curve_Geometry_Result::Segment&& source_segment) {
if (!source_segment.points.empty())
result.segments.push_back(std::move(source_segment));
}
static void merge_curve_first_segment(Curve_Geometry_Result& result, Curve_Geometry_Result::Segment& source_segment) {
if (source_segment.points.empty())
return;
int run_start = -1;
int sample_count = static_cast<int>(samples->values.size());
for (int i = 0; i < sample_count; ++i) {
double value = samples->values[static_cast<std::size_t>(i)];
if (std::isfinite(value)) {
if (run_start < 0)
run_start = i;
if (result.segments.empty()) {
result.segments.push_back(std::move(source_segment));
return;
}
auto& target_points = result.segments.back().points;
for (const PointF& point : source_segment.points) {
if (!target_points.empty() && target_points.back().x == point.x && target_points.back().y == point.y)
continue;
target_points.push_back(point);
}
}
static Curve_Geometry_Result merge_curve_geometry(Curve_Geometry_Job& job) {
Curve_Geometry_Result result;
if (job.request)
result.mapping = job.request->mapping;
bool previous_can_join = false;
for (std::size_t i = 0; i < job.results.size(); ++i) {
Curve_Geometry_Partial_Result& partial = job.results[i];
if (partial.segments.empty()) {
previous_can_join = false;
continue;
}
if (run_start >= 0)
append_curve_sample_run(result, request, *samples, run_start, i - 1);
run_start = -1;
bool join = previous_can_join && partial.joins_previous;
if (join) {
merge_curve_first_segment(result, partial.segments.front());
for (std::size_t s = 1; s < partial.segments.size(); ++s)
append_curve_segment(result, std::move(partial.segments[s]));
}
else {
for (auto& segment : partial.segments)
append_curve_segment(result, std::move(segment));
}
previous_can_join = partial.ends_with_continuous_point;
}
if (run_start >= 0)
append_curve_sample_run(result, request, *samples, run_start, sample_count - 1);
result.update_states = std::move(job.update_states);
return result;
}
struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_Data> {
struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_Data>, Frame_Prepare_Subtask_Source {
static constexpr std::size_t Input_Queue_Capacity = 8;
std::shared_ptr<const Curve_Source> source = std::make_shared<Curve_Source>();
rigtorp::MPMCQueue<std::shared_ptr<Curve_Input_Block>> input_queue;
@@ -303,13 +464,39 @@ struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_
}
void prepare_data(const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) override {
(void)snapshot;
consume_input_queue();
submit_geometry_if_needed(snapshot, frame_view);
(void)frame_view;
}
void submit_geometry_if_needed(const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) {
bool build_prepare_subtasks(Render_Task_Subflow& subflow, const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) override {
consume_input_queue();
auto job = create_geometry_job_if_needed(snapshot, frame_view);
if (!job)
return true;
std::vector<Render_Task_Graph_Node> partition_tasks;
partition_tasks.reserve(job->partitions.size());
for (std::size_t i = 0; i < job->partitions.size(); ++i) {
partition_tasks.push_back(subflow.emplace(Render_Task_Kind::Curve, Task([job, i]() {
build_curve_partition(*job, i);
})));
}
auto merge = subflow.emplace(Render_Task_Kind::Compose, Task([this, job]() {
if (q() && q()->get_retiring())
return;
Curve_Geometry_Result result = merge_curve_geometry(*job);
if (q() && q()->get_retiring()) {
complete_curve_result_states(result, std::make_error_code(std::errc::operation_canceled), Update_Outcome::Cancelled);
return;
}
geometry_result = std::move(result);
submitted_geometry_key = job->key;
}));
for (auto task : partition_tasks)
subflow.precede(task, merge);
return true;
}
std::shared_ptr<Curve_Geometry_Job> create_geometry_job_if_needed(const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) {
Curve_Render_State* s = render_state(frame_view);
if (!s)
return;
return {};
auto domain_axis = s->domain_axis.lock();
auto value_axis = s->value_axis.lock();
if (!domain_axis || !value_axis) {
@@ -318,7 +505,7 @@ struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_
state->complete_all(std::make_error_code(std::errc::operation_canceled), Update_Outcome::Cancelled);
}
pending_geometry_update_states.clear();
return;
return {};
}
Axis_Mapping_2D mapping = Axis_Render_Access::mapping(domain_axis.get(), value_axis.get(), snapshot);
Curve_Geometry_Key key;
@@ -334,8 +521,7 @@ struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_
key.domain = curve_axis_key(mapping.domain);
key.value = curve_axis_key(mapping.value);
if (submitted_geometry_key && key == *submitted_geometry_key)
return;
submitted_geometry_key = key;
return {};
auto request = std::make_shared<Curve_Geometry_Request>();
request->mode = key.mode;
request->mapping = mapping;
@@ -346,10 +532,8 @@ struct Curve_Private : Typed_Render_Data<Curve, Curve_Render_State, Curve_Input_
request->source = source;
request->update_states = std::move(pending_geometry_update_states);
pending_geometry_update_states = std::pmr::vector<std::shared_ptr<Update_State>>(memory_resource(Memory_Domain::Update_Completion));
geometry_result.emplace();
geometry_result->mapping = request->mapping;
build_curve_geometry(*geometry_result, *request);
geometry_result->update_states = std::move(request->update_states);
std::size_t worker_count = std::max<std::size_t>(1, Global::instance()->render_scheduler().render_worker_count());
return create_curve_geometry_job(std::move(request), std::move(key), worker_count);
}
void draw(Canvas& canvas, const Render_Frame_Snapshot& snapshot, const Renderable_Frame_View& frame_view) override {
Curve_Render_State* s = render_state(frame_view);
+4
View File
@@ -0,0 +1,4 @@
#include "Frame_Prepare_Subflow.h"
namespace renderive {
} // namespace renderive
+33
View File
@@ -0,0 +1,33 @@
#pragma once
#include "../architecture/Update_Completion.h"
#include "../execution/Render_Executor.h"
#include <memory>
#include <memory_resource>
namespace renderive {
struct Render_Frame_Snapshot;
struct Renderable_Frame_View;
class Renderable;
class Frame_Prepare_Subtask_Source {
public:
virtual ~Frame_Prepare_Subtask_Source() = default;
virtual bool build_prepare_subtasks(
Render_Task_Subflow& subflow,
const Render_Frame_Snapshot& snapshot,
const Renderable_Frame_View& frame_view) = 0;
};
void build_renderable_prepare(
Renderable& renderable,
Render_Task_Subflow& subflow,
const Render_Frame_Snapshot& snapshot,
const Renderable_Frame_View& frame_view);
void drain_renderable_prepare_updates(
Renderable& renderable,
std::pmr::vector<std::shared_ptr<Update_State>>& output);
} // namespace renderive
+92 -14
View File
@@ -1307,7 +1307,6 @@ TEST(Renderive_Render_Executor, ExecutorDoesNotOwnPlotAdmission) {
completed.fetch_add(1, std::memory_order_acq_rel);
}),
first_stat,
1,
true
}));
ASSERT_TRUE(wait_until([&]() {
@@ -1324,7 +1323,6 @@ TEST(Renderive_Render_Executor, ExecutorDoesNotOwnPlotAdmission) {
completed.fetch_add(1, std::memory_order_acq_rel);
}),
second_stat,
1,
true
}));
auto third_stat = std::make_shared<Frame_Worker_Stats>();
@@ -1338,7 +1336,6 @@ TEST(Renderive_Render_Executor, ExecutorDoesNotOwnPlotAdmission) {
completed.fetch_add(1, std::memory_order_acq_rel);
}),
third_stat,
1,
true
}));
auto duplicate_stat = std::make_shared<Frame_Worker_Stats>();
@@ -1352,7 +1349,6 @@ TEST(Renderive_Render_Executor, ExecutorDoesNotOwnPlotAdmission) {
completed.fetch_add(1, std::memory_order_acq_rel);
}),
duplicate_stat,
1,
true
}));
auto overflow_stat = std::make_shared<Frame_Worker_Stats>();
@@ -1366,7 +1362,6 @@ TEST(Renderive_Render_Executor, ExecutorDoesNotOwnPlotAdmission) {
completed.fetch_add(1, std::memory_order_acq_rel);
}),
overflow_stat,
1,
true
}));
release.store(true, std::memory_order_release);
@@ -1386,15 +1381,15 @@ TEST(Renderive_Render_Executor, TaskflowTopologyRunsInFrameOrder) {
Render_Task_Kind::Frame,
Task{},
[&sequence](Render_Task_Graph& graph, Render_Task_Graph_Node entry, Render_Task_Graph_Node exit) {
auto prepare = graph.emplace(Task([&sequence]() {
auto prepare = graph.emplace(Render_Task_Kind::Frame, Task([&sequence]() {
int expected = 0;
sequence.compare_exchange_strong(expected, 1, std::memory_order_acq_rel);
}));
auto draw = graph.emplace(Task([&sequence]() {
auto draw = graph.emplace(Render_Task_Kind::Primitive, Task([&sequence]() {
int expected = 1;
sequence.compare_exchange_strong(expected, 2, std::memory_order_acq_rel);
}));
auto finalize = graph.emplace(Task([&sequence]() {
auto finalize = graph.emplace(Render_Task_Kind::Frame, Task([&sequence]() {
int expected = 2;
sequence.compare_exchange_strong(expected, 3, std::memory_order_acq_rel);
}));
@@ -1407,7 +1402,6 @@ TEST(Renderive_Render_Executor, TaskflowTopologyRunsInFrameOrder) {
completed.store(true, std::memory_order_release);
}),
stat,
3,
true
}));
ASSERT_TRUE(wait_until([&]() {
@@ -1417,6 +1411,95 @@ TEST(Renderive_Render_Executor, TaskflowTopologyRunsInFrameOrder) {
EXPECT_EQ(stat->task_count, 3u);
executor.shutdown();
}
TEST(Renderive_Render_Executor, FrameDrawWaitsForPrepareSubflow) {
Render_Executor executor(Render_Runtime_Config{2});
std::latch child_entered(1);
std::latch release_child(1);
std::latch completed(1);
std::atomic_bool child_finished{false};
std::atomic_bool draw_started{false};
auto stat = std::make_shared<Frame_Worker_Stats>();
ASSERT_TRUE(executor.try_submit(Render_Executor_Task{
1,
1,
Render_Task_Kind::Frame,
Task{},
[&](Render_Task_Graph& graph, Render_Task_Graph_Node entry, Render_Task_Graph_Node exit) {
auto prepare = graph.emplace_subflow(Render_Task_Kind::Frame, [&](Render_Task_Subflow& subflow) {
auto child = subflow.emplace(Render_Task_Kind::Curve, Task([&]() {
child_entered.count_down();
release_child.wait();
child_finished.store(true, std::memory_order_release);
}));
auto merge = subflow.emplace(Render_Task_Kind::Compose, Task([&]() {
EXPECT_TRUE(child_finished.load(std::memory_order_acquire));
}));
subflow.precede(child, merge);
});
auto draw = graph.emplace(Render_Task_Kind::Primitive, Task([&]() {
draw_started.store(true, std::memory_order_release);
EXPECT_TRUE(child_finished.load(std::memory_order_acquire));
}));
graph.precede(entry, prepare);
graph.precede(prepare, draw);
graph.precede(draw, exit);
},
Task([&]() {
completed.count_down();
}),
stat,
true
}));
child_entered.wait();
EXPECT_FALSE(draw_started.load(std::memory_order_acquire));
release_child.count_down();
completed.wait();
EXPECT_TRUE(draw_started.load(std::memory_order_acquire));
EXPECT_GE(stat->task_count, 3u);
EXPECT_GE(stat->peak_parallelism, 1u);
executor.shutdown();
}
TEST(Renderive_Render_Executor, SubflowKindCountersReflectChildren) {
Render_Executor executor(Render_Runtime_Config{4});
std::latch completed(1);
ASSERT_TRUE(executor.try_submit(Render_Executor_Task{
1,
1,
Render_Task_Kind::Frame,
Task{},
[](Render_Task_Graph& graph, Render_Task_Graph_Node entry, Render_Task_Graph_Node exit) {
auto prepare = graph.emplace_subflow(Render_Task_Kind::Frame, [](Render_Task_Subflow& subflow) {
auto curve0 = subflow.emplace(Render_Task_Kind::Curve, Task([]() {}));
auto curve1 = subflow.emplace(Render_Task_Kind::Curve, Task([]() {}));
auto waterfall0 = subflow.emplace(Render_Task_Kind::Waterfall, Task([]() {}));
auto waterfall1 = subflow.emplace(Render_Task_Kind::Waterfall, Task([]() {}));
auto compose = subflow.emplace(Render_Task_Kind::Compose, Task([]() {}));
subflow.precede(curve0, compose);
subflow.precede(curve1, compose);
subflow.precede(waterfall0, compose);
subflow.precede(waterfall1, compose);
});
auto draw = graph.emplace(Render_Task_Kind::Primitive, Task([]() {}));
auto finalize = graph.emplace(Render_Task_Kind::Frame, Task([]() {}));
graph.precede(entry, prepare);
graph.precede(prepare, draw);
graph.precede(draw, finalize);
graph.precede(finalize, exit);
},
Task([&]() {
completed.count_down();
}),
std::make_shared<Frame_Worker_Stats>(),
true
}));
completed.wait();
executor.shutdown();
Render_Executor_Snapshot snapshot = executor.snapshot();
EXPECT_GE(snapshot.curve_tasks, 2u);
EXPECT_GE(snapshot.waterfall_tasks, 2u);
EXPECT_GE(snapshot.compose_tasks, 1u);
EXPECT_GE(snapshot.primitive_tasks, 3u);
}
TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) {
Render_Executor executor(Render_Runtime_Config{1});
std::atomic_bool first_entered{false};
@@ -1436,7 +1519,6 @@ TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) {
{},
Task([]() {}),
std::make_shared<Frame_Worker_Stats>(),
1,
true
}));
ASSERT_TRUE(wait_until([&]() {
@@ -1452,7 +1534,6 @@ TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) {
{},
Task([]() {}),
std::make_shared<Frame_Worker_Stats>(),
1,
true
}));
ASSERT_TRUE(executor.try_submit(Render_Executor_Task{
@@ -1465,7 +1546,6 @@ TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) {
{},
Task([]() {}),
std::make_shared<Frame_Worker_Stats>(),
1,
true
}));
release_first.store(true, std::memory_order_release);
@@ -1493,7 +1573,6 @@ TEST(Renderive_Render_Executor, ShutdownWaitsForAcceptedQueuedFrame) {
{},
Task([]() {}),
std::make_shared<Frame_Worker_Stats>(),
1,
true
}));
ASSERT_TRUE(wait_until([&]() {
@@ -1511,7 +1590,6 @@ TEST(Renderive_Render_Executor, ShutdownWaitsForAcceptedQueuedFrame) {
queued_completed.store(true, std::memory_order_release);
}),
std::make_shared<Frame_Worker_Stats>(),
1,
true
}));
std::thread shutdown_thread([&]() {