明面上没bug了
This commit is contained in:
@@ -0,0 +1,284 @@
|
||||
#include "Render_Executor.h"
|
||||
#include "../architecture/Render_Time.h"
|
||||
#include <algorithm>
|
||||
#include <array>
|
||||
#include <atomic>
|
||||
#include <memory>
|
||||
#include <mutex>
|
||||
#include <unordered_set>
|
||||
#include <taskflow/taskflow.hpp>
|
||||
#include <thread>
|
||||
|
||||
namespace renderive {
|
||||
namespace {
|
||||
|
||||
struct Render_Executor_Metrics {
|
||||
std::size_t worker_count{};
|
||||
std::atomic<std::size_t> active_workers{0};
|
||||
std::atomic<std::size_t> active_frame_jobs{0};
|
||||
std::atomic<std::size_t> queued_frame_jobs{0};
|
||||
std::atomic<std::uint64_t> submitted_tasks{0};
|
||||
std::atomic<std::uint64_t> started_tasks{0};
|
||||
std::atomic<std::uint64_t> completed_tasks{0};
|
||||
std::atomic<std::uint64_t> rejected_frame_jobs{0};
|
||||
std::atomic<std::uint64_t> dropped_frame_attempts{0};
|
||||
std::atomic<std::uint64_t> peak_concurrency{0};
|
||||
std::atomic<std::uint64_t> queue_wait_total_ns{0};
|
||||
std::atomic<std::uint64_t> queue_wait_max_ns{0};
|
||||
std::atomic<std::uint64_t> task_duration_total_ns{0};
|
||||
std::atomic<std::uint64_t> task_duration_max_ns{0};
|
||||
std::atomic<std::uint64_t> curve_tasks{0};
|
||||
std::atomic<std::uint64_t> waterfall_tasks{0};
|
||||
std::atomic<std::uint64_t> primitive_tasks{0};
|
||||
std::atomic<std::uint64_t> compose_tasks{0};
|
||||
std::atomic<std::uint64_t> performance_tasks{0};
|
||||
};
|
||||
|
||||
void update_peak(std::atomic<std::uint64_t>& 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<Render_Executor_Metrics> 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<std::uint64_t>(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<Render_Executor_Metrics> metrics;
|
||||
std::vector<std::uint64_t> worker_entry_ns;
|
||||
};
|
||||
|
||||
class Render_Executor_Private {
|
||||
public:
|
||||
explicit Render_Executor_Private(std::size_t worker_count)
|
||||
: metrics(std::make_shared<Render_Executor_Metrics>()),
|
||||
executor(worker_count) {
|
||||
metrics->worker_count = worker_count;
|
||||
observer = executor.make_observer<Render_Executor_Observer>(metrics);
|
||||
}
|
||||
|
||||
~Render_Executor_Private() {
|
||||
shutdown();
|
||||
}
|
||||
|
||||
bool try_submit(Render_Executor_Task task) {
|
||||
if (shutting_down.load(std::memory_order_acquire) || !task.work) {
|
||||
reject(task);
|
||||
return false;
|
||||
}
|
||||
if (!admit(task)) {
|
||||
reject(task);
|
||||
return false;
|
||||
}
|
||||
std::uint64_t enqueue_ns = steady_now_ns();
|
||||
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();
|
||||
|
||||
auto topology = std::make_shared<tf::Taskflow>();
|
||||
auto task_ptr = std::make_shared<Render_Executor_Task>(std::move(task));
|
||||
auto completion = std::make_shared<Task>(std::move(task_ptr->completion));
|
||||
topology->emplace([this, enqueue_ns, task_ptr]() mutable {
|
||||
Render_Executor_Task& task = *task_ptr;
|
||||
std::uint64_t begin_ns = steady_now_ns();
|
||||
std::uint64_t queue_wait_ns = begin_ns > enqueue_ns ? begin_ns - enqueue_ns : 0;
|
||||
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);
|
||||
}
|
||||
task.work();
|
||||
std::uint64_t end_ns = steady_now_ns();
|
||||
std::uint64_t run_ns = end_ns > begin_ns ? end_ns - begin_ns : 0;
|
||||
if (task.frame_stat) {
|
||||
Frame_Worker_Stat& stat = *task.frame_stat;
|
||||
stat.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)));
|
||||
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);
|
||||
});
|
||||
executor.run(*topology, [this, topology, completion, task_ptr]() mutable {
|
||||
release(*task_ptr);
|
||||
if (completion && *completion)
|
||||
(*completion)();
|
||||
});
|
||||
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 = 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:
|
||||
bool admit(const Render_Executor_Task& task) {
|
||||
if (!task.frame_job)
|
||||
return true;
|
||||
std::lock_guard<std::mutex> lock(admission_mutex);
|
||||
std::size_t active = metrics->active_frame_jobs.load(std::memory_order_acquire);
|
||||
std::size_t queued = metrics->queued_frame_jobs.load(std::memory_order_acquire);
|
||||
std::size_t limit = std::max<std::size_t>(1, metrics->worker_count);
|
||||
if (active + queued >= limit)
|
||||
return false;
|
||||
if (task.plot_id && !admitted_frame_plots.insert(task.plot_id).second)
|
||||
return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
void release(const Render_Executor_Task& task) {
|
||||
if (!task.frame_job || !task.plot_id)
|
||||
return;
|
||||
std::lock_guard<std::mutex> lock(admission_mutex);
|
||||
admitted_frame_plots.erase(task.plot_id);
|
||||
}
|
||||
|
||||
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<Render_Executor_Metrics> metrics;
|
||||
tf::Executor executor;
|
||||
std::shared_ptr<tf::ObserverInterface> observer;
|
||||
std::mutex admission_mutex;
|
||||
std::unordered_set<Plot_Execution_Id> admitted_frame_plots;
|
||||
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<Render_Executor_Private>(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
|
||||
Reference in New Issue
Block a user