#include "Render_Executor.h" #include "../architecture/Render_Time.h" #include #include #include #include #include #include #include #include #include #include #include namespace renderive { namespace { 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_tasks{0}; std::atomic started_tasks{0}; std::atomic completed_tasks{0}; std::atomic rejected_frame_jobs{0}; std::atomic dropped_frame_attempts{0}; std::atomic peak_concurrency{0}; std::atomic queue_wait_total_ns{0}; std::atomic queue_wait_max_ns{0}; std::atomic task_duration_total_ns{0}; std::atomic task_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; } } } // namespace 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_tasks.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_duration_total_ns.fetch_add(duration_ns, std::memory_order_relaxed); update_peak(metrics->task_duration_max_ns, duration_ns); } metrics->completed_tasks.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 Task_Run_Record { std::uint64_t begin_ns{}; std::uint64_t queue_wait_ns{}; }; struct Pending_Frame_Task { Render_Executor_Task task; std::uint64_t enqueue_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_taskflow)) { reject(task); return false; } std::uint64_t enqueue_ns = steady_now_ns(); if (!task.frame_job) { record_submission(task); submit_to_executor(std::move(task), enqueue_ns); return true; } bool submit_now = false; { std::lock_guard lock(admission_mutex); if (plot_already_pending_or_admitted_locked(task.plot_id)) { reject(task); return false; } std::size_t limit = frame_admission_limit(); if (admitted_frame_count < limit) { admit_frame_locked(task.plot_id); submit_now = true; } else { std::size_t queue_limit = frame_queue_limit(); if (pending_frame_tasks.size() >= queue_limit) { reject(task); return false; } if (task.plot_id) queued_frame_plots.insert(task.plot_id); record_submission(task); pending_frame_tasks.push_back(Pending_Frame_Task{std::move(task), enqueue_ns}); return true; } } if (submit_now) { record_submission(task); submit_to_executor(std::move(task), enqueue_ns); return true; } reject(task); return false; } void shutdown() { if (shutting_down.exchange(true, std::memory_order_acq_rel)) return; std::vector cancelled_tasks; { std::lock_guard lock(admission_mutex); while (!pending_frame_tasks.empty()) { Pending_Frame_Task pending = std::move(pending_frame_tasks.front()); pending_frame_tasks.pop_front(); cancelled_tasks.push_back(std::move(pending.task)); metrics->queued_frame_jobs.fetch_sub(1, std::memory_order_relaxed); } pending_frame_tasks.clear(); queued_frame_plots.clear(); } for (auto& task : cancelled_tasks) { if (task.frame_stat) task.frame_stat->executor_at_finish = snapshot(); if (task.completion) task.completion(); } executor.wait_for_all(); } Render_Executor_Snapshot snapshot() const { std::uint64_t submitted = metrics->submitted_tasks.load(std::memory_order_acquire); std::uint64_t completed = metrics->completed_tasks.load(std::memory_order_acquire); std::uint64_t started = metrics->started_tasks.load(std::memory_order_acquire); std::uint64_t queue_total = metrics->queue_wait_total_ns.load(std::memory_order_acquire); std::uint64_t duration_total = metrics->task_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, started, completed, 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), submitted ? queue_total / submitted : 0, metrics->queue_wait_max_ns.load(std::memory_order_acquire), completed ? duration_total / completed : 0, metrics->task_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_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) 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 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->queue_wait_total_ns.fetch_add(queue_wait_ns, std::memory_order_relaxed); update_peak(metrics->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]() 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; 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))); 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.parallel_stage_wall_ns += run_ns; stat.executor_at_finish = snapshot(); } if (task.frame_job) metrics->active_frame_jobs.fetch_sub(1, std::memory_order_relaxed); }); if (task_ptr->build_taskflow) { task_ptr->build_taskflow(*topology, entry, exit); } else { auto work = topology->emplace([task_ptr]() mutable { if (task_ptr->work) task_ptr->work(); }); entry.precede(work); work.precede(exit); } executor.run(*topology, [this, topology, completion, task_ptr]() mutable { release_and_drain_next(*task_ptr); if (completion && *completion) (*completion)(); }); } std::size_t frame_admission_limit() const { return std::max(1, metrics->worker_count); } std::size_t frame_queue_limit() const { return std::max(1, frame_admission_limit() * 2); } bool plot_already_pending_or_admitted_locked(Plot_Execution_Id plot_id) const { if (!plot_id) return false; return admitted_frame_plots.find(plot_id) != admitted_frame_plots.end() || queued_frame_plots.find(plot_id) != queued_frame_plots.end(); } void admit_frame_locked(Plot_Execution_Id plot_id) { ++admitted_frame_count; if (plot_id) admitted_frame_plots.insert(plot_id); } std::optional take_next_frame_task_locked() { if (shutting_down.load(std::memory_order_acquire) || pending_frame_tasks.empty()) return std::nullopt; if (admitted_frame_count >= frame_admission_limit()) return std::nullopt; Pending_Frame_Task next = std::move(pending_frame_tasks.front()); pending_frame_tasks.pop_front(); if (next.task.plot_id) queued_frame_plots.erase(next.task.plot_id); admit_frame_locked(next.task.plot_id); return next; } void release_and_drain_next(const Render_Executor_Task& task) { if (!task.frame_job) return; std::optional next; { std::lock_guard lock(admission_mutex); if (admitted_frame_count) --admitted_frame_count; if (task.plot_id) admitted_frame_plots.erase(task.plot_id); next = take_next_frame_task_locked(); } if (next) submit_to_executor(std::move(next->task), next->enqueue_ns); } 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_task_count++; } std::shared_ptr metrics; tf::Executor executor; std::shared_ptr observer; std::mutex admission_mutex; std::unordered_set admitted_frame_plots; std::unordered_set queued_frame_plots; std::deque pending_frame_tasks; std::size_t admitted_frame_count{}; 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