diff --git a/Core/architecture/Renderable.cpp b/Core/architecture/Renderable.cpp index 013a8ee..2d8f631 100644 --- a/Core/architecture/Renderable.cpp +++ b/Core/architecture/Renderable.cpp @@ -1,6 +1,7 @@ #include #include #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(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>& 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; diff --git a/Core/architecture/Renderable.h b/Core/architecture/Renderable.h index 1c1b99b..1a0b75b 100644 --- a/Core/architecture/Renderable.h +++ b/Core/architecture/Renderable.h @@ -24,6 +24,7 @@ namespace renderive { class Plot_Core; class Renderable; +class Render_Task_Subflow; struct Plot_Render_Context; template 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>& output); template friend struct Renderable_Binding; template diff --git a/Core/execution/Render_Executor.cpp b/Core/execution/Render_Executor.cpp index 0749c7f..252b879 100644 --- a/Core/execution/Render_Executor.cpp +++ b/Core/execution/Render_Executor.cpp @@ -9,37 +9,42 @@ #include 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 metrics, + std::shared_ptr frame_counters) + : taskflow(taskflow), + metrics(std::move(metrics)), + frame_counters(std::move(frame_counters)) {} tf::Taskflow& taskflow; + std::shared_ptr metrics; + std::shared_ptr frame_counters; std::vector 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(impl); - if (!graph || !work) - return {}; - auto work_ptr = std::make_shared(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(graph->tasks.size() - 1)}; -} - -void Render_Task_Graph::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) { - auto* graph = static_cast(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 metrics, + std::shared_ptr frame_counters) + : subflow(subflow), + metrics(std::move(metrics)), + frame_counters(std::move(frame_counters)) {} + tf::Subflow& subflow; + std::shared_ptr metrics; + std::shared_ptr frame_counters; + std::vector 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 metrics, + std::shared_ptr 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(impl); + if (!graph || !work) + return {}; + auto wrapped = std::make_shared(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(graph->tasks.size() - 1)}; +} + +void Render_Task_Subflow::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) { + auto* graph = static_cast(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(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(impl); + if (!graph || !work) + return {}; + auto wrapped = std::make_shared(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(graph->tasks.size() - 1)}; +} + +Render_Task_Graph_Node Render_Task_Graph::emplace_subflow(Render_Task_Kind kind, Subflow_Builder builder) { + auto* graph = static_cast(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(graph->tasks.size() - 1)}; +} + +void Render_Task_Graph::precede(Render_Task_Graph_Node before, Render_Task_Graph_Node after) { + auto* graph = static_cast(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(std::move(task)); auto completion = std::make_shared(std::move(task_ptr->completion)); auto run_record = std::make_shared(); + auto frame_counters = std::make_shared(); auto topology = std::make_shared(); 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(1, task.logical_task_count); - stat.peak_parallelism = static_cast(std::max(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(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); diff --git a/Core/execution/Render_Executor.h b/Core/execution/Render_Executor.h index 72d7016..edaf77a 100644 --- a/Core/execution/Render_Executor.h +++ b/Core/execution/Render_Executor.h @@ -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; 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_stat; - std::uint32_t logical_task_count{1}; bool frame_job{}; }; class Render_Executor_Private; diff --git a/Core/plot/Plot_Core.cpp b/Core/plot/Plot_Core.cpp index 320a183..67fdf24 100644 --- a/Core/plot/Plot_Core.cpp +++ b/Core/plot/Plot_Core.cpp @@ -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 #include @@ -21,6 +22,7 @@ struct Frame_Render_Job { Render_Frame_Snapshot render_snapshot; Renderable_Frame_Tree renderables; std::shared_ptr worker_stat; + std::atomic_uint64_t prepared_output_count{0}; }; class Plot_Core_Private final : public std::enable_shared_from_this, 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& job) { return job && job->buffer && job->lifecycle && !job->render_snapshot.destroying(); } - void prepare_frame_job(const std::shared_ptr& job) { + void begin_prepare_frame_job(const std::shared_ptr& 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& 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& 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& 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); } diff --git a/Core/plottable/Waterfall_p.h b/Core/plottable/Waterfall_p.h index 4fac452..4358de9 100644 --- a/Core/plottable/Waterfall_p.h +++ b/Core/plottable/Waterfall_p.h @@ -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 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 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(columns * rows)); } - std::shared_ptr 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 request; + std::pmr::vector> update_states; int tile_columns{}; int tile_rows{}; std::vector tiles; - std::atomic 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 create_waterfall_image_job(std::shared_ptr 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> update_states = std::move(request->update_states); + auto job = std::make_shared( + std::shared_ptr(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(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& 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(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, Hit_Testable, Hover_Interactive { +struct Waterfall_Private : Typed_Render_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; @@ -357,6 +371,13 @@ struct Waterfall_Private : Typed_Render_Datafrequency_axis.lock(); @@ -469,7 +490,8 @@ struct Waterfall_Private : Typed_Render_Data(); request->key = key; request->color_map = color_scale->color_map; @@ -497,7 +519,29 @@ struct Waterfall_Private : Typed_Render_Dataupdate_states = std::move(pending_image_update_states); pending_image_update_states = std::pmr::vector>(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 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(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) { diff --git a/Core/primitive/Curve.cpp b/Core/primitive/Curve.cpp index 91c9c4d..6555f54 100644 --- a/Core/primitive/Curve.cpp +++ b/Core/primitive/Curve.cpp @@ -11,7 +11,9 @@ #include #include #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 source; std::pmr::vector> 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 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 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 request; + Curve_Geometry_Key key; + std::vector partitions; + std::vector results; + std::pmr::vector> 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(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(request.source.get())) + return data_points->values.size(); + if (const auto* samples = std::get_if(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(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(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(samples->values.size()), static_cast(index)); + return finite_point(request.mapping.map(coord, value)); + } + return false; +} + +static std::shared_ptr create_curve_geometry_job(std::shared_ptr 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> update_states = std::move(request->update_states); + auto job = std::make_shared( + std::shared_ptr(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(request.source.get()); - if (!samples || samples->values.size() < 2 || request.data_range.length() == 0.0) + if (const auto* data_points = std::get_if(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(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(i); + continue; + } + if (run_start >= 0) + append_curve_sample_run(result, request, *samples, run_start, static_cast(i) - 1); + run_start = -1; + } + if (run_start >= 0) + append_curve_sample_run(result, request, *samples, run_start, static_cast(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(samples->values.size()); - for (int i = 0; i < sample_count; ++i) { - double value = samples->values[static_cast(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 { +struct Curve_Private : Typed_Render_Data, Frame_Prepare_Subtask_Source { static constexpr std::size_t Input_Queue_Capacity = 8; std::shared_ptr source = std::make_shared(); rigtorp::MPMCQueue> input_queue; @@ -303,13 +464,39 @@ struct Curve_Private : Typed_Render_Data 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 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_Datacomplete_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(); request->mode = key.mode; request->mapping = mapping; @@ -346,10 +532,8 @@ struct Curve_Private : Typed_Render_Datasource = source; request->update_states = std::move(pending_geometry_update_states); pending_geometry_update_states = std::pmr::vector>(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(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); diff --git a/Core/render/Frame_Prepare_Subflow.cpp b/Core/render/Frame_Prepare_Subflow.cpp new file mode 100644 index 0000000..4dd272c --- /dev/null +++ b/Core/render/Frame_Prepare_Subflow.cpp @@ -0,0 +1,4 @@ +#include "Frame_Prepare_Subflow.h" + +namespace renderive { +} // namespace renderive diff --git a/Core/render/Frame_Prepare_Subflow.h b/Core/render/Frame_Prepare_Subflow.h new file mode 100644 index 0000000..3fc77f5 --- /dev/null +++ b/Core/render/Frame_Prepare_Subflow.h @@ -0,0 +1,33 @@ +#pragma once +#include "../architecture/Update_Completion.h" +#include "../execution/Render_Executor.h" +#include +#include + +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>& output); + +} // namespace renderive diff --git a/test/Renderive_Core_Tests.cpp b/test/Renderive_Core_Tests.cpp index 2f86bad..77053ab 100644 --- a/test/Renderive_Core_Tests.cpp +++ b/test/Renderive_Core_Tests.cpp @@ -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(); @@ -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(); @@ -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(); @@ -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(); + 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(), + 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(), - 1, true })); ASSERT_TRUE(wait_until([&]() { @@ -1452,7 +1534,6 @@ TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) { {}, Task([]() {}), std::make_shared(), - 1, true })); ASSERT_TRUE(executor.try_submit(Render_Executor_Task{ @@ -1465,7 +1546,6 @@ TEST(Renderive_Render_Executor, SubmittedFramesAllCompleteOnSingleWorker) { {}, Task([]() {}), std::make_shared(), - 1, true })); release_first.store(true, std::memory_order_release); @@ -1493,7 +1573,6 @@ TEST(Renderive_Render_Executor, ShutdownWaitsForAcceptedQueuedFrame) { {}, Task([]() {}), std::make_shared(), - 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(), - 1, true })); std::thread shutdown_thread([&]() {