#include "Render_Executor.h" #include "../architecture/Render_Time.h" #include #include #include #include #include #include #include namespace renderive { struct Render_Executor_Metrics; struct Frame_Business_Task_Counters { std::atomic_uint32_t business_task_count{0}; std::atomic_uint32_t active_business_task_count{0}; std::atomic_uint32_t business_task_peak_parallelism{0}; std::atomic_uint64_t business_task_duration_total_ns{0}; }; struct Taskflow_Graph_Impl { 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; }; 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{}; std::atomic active_workers{0}; std::atomic active_frame_jobs{0}; std::atomic queued_frame_jobs{0}; std::atomic submitted_topologies{0}; std::atomic started_topologies{0}; std::atomic completed_topologies{0}; std::atomic started_task_nodes{0}; std::atomic completed_task_nodes{0}; std::atomic rejected_frame_jobs{0}; std::atomic dropped_frame_attempts{0}; std::atomic peak_concurrency{0}; std::atomic topology_queue_wait_total_ns{0}; std::atomic topology_queue_wait_max_ns{0}; std::atomic task_node_duration_total_ns{0}; std::atomic task_node_duration_max_ns{0}; std::atomic curve_tasks{0}; std::atomic waterfall_tasks{0}; std::atomic primitive_tasks{0}; std::atomic compose_tasks{0}; std::atomic performance_tasks{0}; }; void update_peak(std::atomic& target, std::uint64_t value) { std::uint64_t current = target.load(std::memory_order_acquire); while (current < value && !target.compare_exchange_weak(current, value, std::memory_order_acq_rel, std::memory_order_acquire)) {} } void add_kind_count(Render_Executor_Metrics& metrics, Render_Task_Kind kind) { switch (kind) { case Render_Task_Kind::Frame: case Render_Task_Kind::Primitive: metrics.primitive_tasks.fetch_add(1, std::memory_order_relaxed); break; case Render_Task_Kind::Curve: metrics.curve_tasks.fetch_add(1, std::memory_order_relaxed); break; case Render_Task_Kind::Waterfall: metrics.waterfall_tasks.fetch_add(1, std::memory_order_relaxed); break; case Render_Task_Kind::Compose: metrics.compose_tasks.fetch_add(1, std::memory_order_relaxed); break; case Render_Task_Kind::Performance: metrics.performance_tasks.fetch_add(1, std::memory_order_relaxed); break; } } void begin_business_task(Render_Executor_Metrics* metrics, Frame_Business_Task_Counters* counters, Render_Task_Kind kind) { if (metrics) add_kind_count(*metrics, kind); if (!counters) return; counters->business_task_count.fetch_add(1, std::memory_order_relaxed); std::uint32_t active = counters->active_business_task_count.fetch_add(1, std::memory_order_acq_rel) + 1; std::uint32_t peak = counters->business_task_peak_parallelism.load(std::memory_order_acquire); while (peak < active && !counters->business_task_peak_parallelism.compare_exchange_weak(peak, active, std::memory_order_acq_rel, std::memory_order_acquire)) {} } void end_business_task(Frame_Business_Task_Counters* counters, std::uint64_t begin_ns) { if (!counters) return; std::uint64_t end_ns = steady_now_ns(); counters->business_task_duration_total_ns.fetch_add(end_ns > begin_ns ? end_ns - begin_ns : 0, std::memory_order_relaxed); counters->active_business_task_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: explicit Render_Executor_Observer(std::shared_ptr metrics) : metrics(std::move(metrics)) {} void set_up(std::size_t worker_count) override { if (metrics) metrics->worker_count = worker_count; worker_entry_ns.clear(); worker_entry_ns.resize(worker_count); } void on_entry(tf::WorkerView worker, tf::TaskView) override { if (!metrics) return; std::size_t worker_id = worker.id(); if (worker_id < worker_entry_ns.size()) worker_entry_ns[worker_id] = steady_now_ns(); std::size_t active = metrics->active_workers.fetch_add(1, std::memory_order_acq_rel) + 1; metrics->started_task_nodes.fetch_add(1, std::memory_order_relaxed); update_peak(metrics->peak_concurrency, static_cast(active)); } void on_exit(tf::WorkerView worker, tf::TaskView) override { if (!metrics) return; std::size_t worker_id = worker.id(); if (worker_id < worker_entry_ns.size()) { std::uint64_t begin_ns = worker_entry_ns[worker_id]; std::uint64_t end_ns = steady_now_ns(); std::uint64_t duration_ns = end_ns > begin_ns ? end_ns - begin_ns : 0; metrics->task_node_duration_total_ns.fetch_add(duration_ns, std::memory_order_relaxed); update_peak(metrics->task_node_duration_max_ns, duration_ns); } metrics->completed_task_nodes.fetch_add(1, std::memory_order_relaxed); metrics->active_workers.fetch_sub(1, std::memory_order_acq_rel); } private: std::shared_ptr metrics; std::vector worker_entry_ns; }; class Render_Executor_Private { struct Topology_Run_Record { std::uint64_t begin_ns{}; std::uint64_t queue_wait_ns{}; }; public: explicit Render_Executor_Private(std::size_t worker_count) : metrics(std::make_shared()), executor(worker_count) { metrics->worker_count = worker_count; observer = executor.make_observer(metrics); } ~Render_Executor_Private() { shutdown(); } bool try_submit(Render_Executor_Task task) { if (shutting_down.load(std::memory_order_acquire) || (!task.work && !task.build_graph)) { reject(task); return false; } std::uint64_t enqueue_ns = steady_now_ns(); record_submission(task); submit_to_executor(std::move(task), enqueue_ns); return true; } void shutdown() { if (shutting_down.exchange(true, std::memory_order_acq_rel)) return; executor.wait_for_all(); } Render_Executor_Snapshot snapshot() const { std::uint64_t submitted_topologies = metrics->submitted_topologies.load(std::memory_order_acquire); std::uint64_t started_topologies = metrics->started_topologies.load(std::memory_order_acquire); std::uint64_t completed_topologies = metrics->completed_topologies.load(std::memory_order_acquire); std::uint64_t completed_task_nodes = metrics->completed_task_nodes.load(std::memory_order_acquire); std::uint64_t started_task_nodes = metrics->started_task_nodes.load(std::memory_order_acquire); std::uint64_t queue_total = metrics->topology_queue_wait_total_ns.load(std::memory_order_acquire); std::uint64_t duration_total = metrics->task_node_duration_total_ns.load(std::memory_order_acquire); std::size_t active = metrics->active_workers.load(std::memory_order_acquire); std::size_t workers = metrics->worker_count; return { workers, active, workers > active ? workers - active : 0, metrics->active_frame_jobs.load(std::memory_order_acquire), metrics->queued_frame_jobs.load(std::memory_order_acquire), submitted_topologies, started_topologies, completed_topologies, started_task_nodes, completed_task_nodes, metrics->rejected_frame_jobs.load(std::memory_order_acquire), metrics->dropped_frame_attempts.load(std::memory_order_acquire), metrics->peak_concurrency.load(std::memory_order_acquire), started_topologies ? queue_total / started_topologies : 0, metrics->topology_queue_wait_max_ns.load(std::memory_order_acquire), completed_task_nodes ? duration_total / completed_task_nodes : 0, metrics->task_node_duration_max_ns.load(std::memory_order_acquire), metrics->curve_tasks.load(std::memory_order_acquire), metrics->waterfall_tasks.load(std::memory_order_acquire), metrics->primitive_tasks.load(std::memory_order_acquire), metrics->compose_tasks.load(std::memory_order_acquire), metrics->performance_tasks.load(std::memory_order_acquire) }; } std::size_t worker_count() const { return metrics->worker_count; } private: void record_submission(Render_Executor_Task& task) { metrics->submitted_topologies.fetch_add(1, std::memory_order_relaxed); if (task.frame_job) metrics->queued_frame_jobs.fetch_add(1, std::memory_order_relaxed); if (task.frame_stat) task.frame_stat->executor_at_submit = snapshot(); } void submit_to_executor(Render_Executor_Task task, std::uint64_t enqueue_ns) { 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; std::uint64_t begin_ns = steady_now_ns(); run_record->begin_ns = begin_ns; std::uint64_t queue_wait_ns = begin_ns > enqueue_ns ? begin_ns - enqueue_ns : 0; run_record->queue_wait_ns = queue_wait_ns; metrics->started_topologies.fetch_add(1, std::memory_order_relaxed); metrics->topology_queue_wait_total_ns.fetch_add(queue_wait_ns, std::memory_order_relaxed); update_peak(metrics->topology_queue_wait_max_ns, queue_wait_ns); if (task.frame_job) { metrics->queued_frame_jobs.fetch_sub(1, std::memory_order_relaxed); metrics->active_frame_jobs.fetch_add(1, std::memory_order_relaxed); } }); 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; std::uint64_t queue_wait_ns = run_record->queue_wait_ns; 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; std::uint32_t business_task_count = frame_counters ? frame_counters->business_task_count.load(std::memory_order_acquire) : 0; std::uint32_t business_task_peak_parallelism = frame_counters ? frame_counters->business_task_peak_parallelism.load(std::memory_order_acquire) : 0; std::uint64_t business_task_duration_total_ns = frame_counters ? frame_counters->business_task_duration_total_ns.load(std::memory_order_acquire) : 0; stat.business_task_count += business_task_count; stat.business_task_peak_parallelism = std::max(stat.business_task_peak_parallelism, business_task_peak_parallelism); stat.topology_queue_wait_total_ns += queue_wait_ns; stat.topology_queue_wait_max_ns = std::max(stat.topology_queue_wait_max_ns, queue_wait_ns); stat.business_task_run_total_ns += business_task_duration_total_ns; stat.topology_run_total_ns += run_ns; } if (task.frame_job) metrics->active_frame_jobs.fetch_sub(1, std::memory_order_relaxed); metrics->completed_topologies.fetch_add(1, std::memory_order_relaxed); }); if (task_ptr->build_graph) { 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 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); } executor.run(*topology, [this, topology, completion, task_ptr]() mutable { if (task_ptr->frame_stat) task_ptr->frame_stat->executor_at_finish = snapshot(); if (completion && *completion) (*completion)(); }); } void reject(const Render_Executor_Task& task) { if (task.frame_job) metrics->rejected_frame_jobs.fetch_add(1, std::memory_order_relaxed); else metrics->dropped_frame_attempts.fetch_add(1, std::memory_order_relaxed); if (task.frame_stat) task.frame_stat->rejected_topology_count++; } std::shared_ptr metrics; tf::Executor executor; std::shared_ptr observer; std::atomic_bool shutting_down{false}; }; std::size_t resolve_render_worker_count(std::size_t configured_count) { if (configured_count) return configured_count; auto cpu_count = std::thread::hardware_concurrency(); return cpu_count > 1 ? cpu_count - 1 : 1; } Render_Executor::Render_Executor(Render_Runtime_Config config) : d(std::make_unique(resolve_render_worker_count(config.worker_count))) {} Render_Executor::~Render_Executor() = default; bool Render_Executor::try_submit(Render_Executor_Task task) { return d->try_submit(std::move(task)); } void Render_Executor::shutdown() { d->shutdown(); } Render_Executor_Snapshot Render_Executor::snapshot() const { return d->snapshot(); } std::size_t Render_Executor::worker_count() const { return d->worker_count(); } } // namespace renderive