250 lines
10 KiB
C++
250 lines
10 KiB
C++
#include "Render_Domain.hpp"
|
|
#include <mutex>
|
|
#include <stdexcept>
|
|
#include <unordered_map>
|
|
namespace aethera::render_3d::detail {
|
|
namespace {
|
|
struct Render_Domain_Registry {
|
|
std::mutex mutex;
|
|
std::unordered_map<std::uint32_t, std::weak_ptr<Render_Domain>> domains;
|
|
};
|
|
Render_Domain_Registry& registry() {
|
|
static Render_Domain_Registry value;
|
|
return value;
|
|
}
|
|
}
|
|
Render_Domain::Prepared_Task::~Prepared_Task() {
|
|
release();
|
|
}
|
|
Render_Domain::Prepared_Task::Prepared_Task(Prepared_Task&& other) noexcept : domain_(std::move(other.domain_)), task_(std::move(other.task_)) {}
|
|
Render_Domain::Prepared_Task& Render_Domain::Prepared_Task::operator=(Prepared_Task&& other) noexcept {
|
|
if (this == &other) return *this;
|
|
release();
|
|
domain_ = std::move(other.domain_);
|
|
task_ = std::move(other.task_);
|
|
return *this;
|
|
}
|
|
void Render_Domain::Prepared_Task::release() noexcept {
|
|
if (!domain_) return;
|
|
task_.reset();
|
|
domain_->release_admission();
|
|
domain_.reset();
|
|
}
|
|
std::shared_ptr<Render_Domain> Render_Domain::acquire(std::uint32_t gpu_index) {
|
|
auto& storage = registry();
|
|
std::lock_guard lock(storage.mutex);
|
|
auto& entry = storage.domains[gpu_index];
|
|
if (auto domain = entry.lock()) return domain;
|
|
auto domain = std::shared_ptr<Render_Domain>(new Render_Domain(),
|
|
&Render_Domain::destroy);
|
|
entry = domain;
|
|
return domain;
|
|
}
|
|
Render_Domain::Render_Domain() {
|
|
thread_ = std::thread([this] {
|
|
run();
|
|
});
|
|
}
|
|
Render_Domain::~Render_Domain() {
|
|
request_stop();
|
|
if (thread_.joinable()) thread_.join();
|
|
}
|
|
void Render_Domain::destroy(Render_Domain* domain) noexcept {
|
|
if (!domain) return;
|
|
if (current_domain_ != domain) {
|
|
delete domain;
|
|
return;
|
|
}
|
|
domain->stopping_.store(true, std::memory_order_release);
|
|
domain->destroy_on_exit_.store(true, std::memory_order_release);
|
|
}
|
|
void Render_Domain::request_stop() noexcept {
|
|
bool expected = false;
|
|
if (!stopping_.compare_exchange_strong(
|
|
expected, true, std::memory_order_acq_rel))
|
|
return;
|
|
slots_.acquire();
|
|
{ std::lock_guard lock(task_mutex_); tasks_.push_back({}); }
|
|
task_condition_.notify_one();
|
|
}
|
|
void Render_Domain::update_peak(std::atomic_size_t& peak, std::size_t value) noexcept {
|
|
std::size_t current = peak.load(std::memory_order_relaxed);
|
|
while (current < value && !peak.compare_exchange_weak(current, value, std::memory_order_relaxed)) {}
|
|
}
|
|
void Render_Domain::acquire_admission() {
|
|
const auto started = std::chrono::steady_clock::now();
|
|
if (!slots_.try_acquire()) {
|
|
backpressure_count_.fetch_add(1, std::memory_order_relaxed);
|
|
slots_.acquire();
|
|
const auto waited = std::chrono::duration_cast<std::chrono::nanoseconds>(
|
|
std::chrono::steady_clock::now() - started).count();
|
|
if (waited > 0) backpressure_wait_ns_.fetch_add(static_cast<std::uint64_t>(waited), std::memory_order_relaxed);
|
|
}
|
|
const std::size_t admitted = admitted_.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_peak(peak_admitted_, admitted);
|
|
}
|
|
bool Render_Domain::try_acquire_admission() noexcept {
|
|
if (!slots_.try_acquire()) {
|
|
backpressure_count_.fetch_add(1, std::memory_order_relaxed);
|
|
return false;
|
|
}
|
|
const std::size_t admitted = admitted_.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_peak(peak_admitted_, admitted);
|
|
return true;
|
|
}
|
|
void Render_Domain::release_admission() noexcept {
|
|
admitted_.fetch_sub(1, std::memory_order_relaxed);
|
|
slots_.release();
|
|
}
|
|
Render_Domain::Prepare_Result Render_Domain::prepare(
|
|
std::function<void()> function, Exception_Handler on_exception) {
|
|
if (!function) throw std::invalid_argument("render domain task is empty");
|
|
if (!on_exception) throw std::invalid_argument("render domain exception handler is empty");
|
|
if (stopping_.load(std::memory_order_acquire)) return {{}, Admission_Result::stopping};
|
|
acquire_admission();
|
|
if (stopping_.load(std::memory_order_acquire)) {
|
|
release_admission();
|
|
return {{}, Admission_Result::stopping};
|
|
}
|
|
try {
|
|
return {
|
|
Prepared_Task(shared_from_this(), std::make_unique<Task>(
|
|
std::move(function), std::move(on_exception))),
|
|
Admission_Result::none
|
|
};
|
|
}
|
|
catch (...) {
|
|
release_admission();
|
|
raise_context("preparing render domain task", std::current_exception());
|
|
}
|
|
}
|
|
Render_Domain::Admission_Result Render_Domain::post(
|
|
std::function<void()> function, Exception_Handler on_exception) {
|
|
auto prepared = prepare(std::move(function), std::move(on_exception));
|
|
if (!prepared) return prepared.result;
|
|
return post(std::move(prepared.task));
|
|
}
|
|
Render_Domain::Try_Post_Result Render_Domain::try_post(std::function<void()> function, Exception_Handler on_exception) {
|
|
if (!function) throw std::invalid_argument("render domain task is empty");
|
|
if (!on_exception) throw std::invalid_argument("render domain exception handler is empty");
|
|
if (stopping_.load(std::memory_order_acquire)) return Try_Post_Result::stopping;
|
|
if (!try_acquire_admission()) return Try_Post_Result::queue_full;
|
|
if (stopping_.load(std::memory_order_acquire)) { release_admission(); return Try_Post_Result::stopping; }
|
|
try {
|
|
auto task = Prepared_Task(shared_from_this(), std::make_unique<Task>(std::move(function), std::move(on_exception)));
|
|
{ std::lock_guard lock(task_mutex_); tasks_.push_back(std::move(task.task_)); }
|
|
task_condition_.notify_one();
|
|
const std::size_t queued = queued_.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_peak(peak_queued_, queued);
|
|
task.domain_.reset();
|
|
return Try_Post_Result::queued;
|
|
}
|
|
catch (...) {
|
|
release_admission();
|
|
raise_context("trying to post render domain task", std::current_exception());
|
|
}
|
|
}
|
|
Render_Domain::Try_Post_Result Render_Domain::try_post_frame(Frame_Task function, Exception_Handler on_exception) {
|
|
if (!function) throw std::invalid_argument("render domain frame task is empty");
|
|
if (!on_exception) throw std::invalid_argument("render domain frame exception handler is empty");
|
|
if (stopping_.load(std::memory_order_acquire)) return Try_Post_Result::stopping;
|
|
auto task = std::make_shared<Frame_Task_Entry>(Frame_Task_Entry{std::move(function), std::move(on_exception)});
|
|
{
|
|
std::lock_guard lock(frame_task_mutex_);
|
|
if (frame_task_in_flight_) {
|
|
if (frame_tasks_.size() >= static_cast<std::size_t>(default_capacity - 1)) return Try_Post_Result::queue_full;
|
|
frame_tasks_.push_back(std::move(task));
|
|
return Try_Post_Result::queued;
|
|
}
|
|
frame_task_in_flight_ = true;
|
|
}
|
|
const auto queued = try_post([this, task] { run_frame_task(task); }, [this, task](std::exception_ptr exception) {
|
|
task->on_exception(std::move(exception));
|
|
complete_frame_task();
|
|
});
|
|
if (queued != Try_Post_Result::queued) {
|
|
std::lock_guard lock(frame_task_mutex_);
|
|
frame_task_in_flight_ = false;
|
|
}
|
|
return queued;
|
|
}
|
|
void Render_Domain::run_frame_task(std::shared_ptr<Frame_Task_Entry> task) {
|
|
auto completed = std::make_shared<std::atomic_bool>();
|
|
auto completion = [this, completed] {
|
|
if (completed->exchange(true, std::memory_order_acq_rel)) return;
|
|
complete_frame_task();
|
|
};
|
|
try {
|
|
std::invoke(task->function, completion);
|
|
}
|
|
catch (...) {
|
|
task->on_exception(contextual_exception("executing render domain frame task", std::current_exception()));
|
|
completion();
|
|
}
|
|
}
|
|
void Render_Domain::complete_frame_task() {
|
|
if (current_domain_ != this) throw std::logic_error("render domain frame completion must run on its domain");
|
|
std::shared_ptr<Frame_Task_Entry> next;
|
|
{
|
|
std::lock_guard lock(frame_task_mutex_);
|
|
if (frame_tasks_.empty()) {
|
|
frame_task_in_flight_ = false;
|
|
return;
|
|
}
|
|
next = std::move(frame_tasks_.front());
|
|
frame_tasks_.pop_front();
|
|
}
|
|
run_frame_task(std::move(next));
|
|
}
|
|
Render_Domain::Admission_Result Render_Domain::post(Prepared_Task task) {
|
|
if (!task.task_ || task.domain_.get() != this) throw std::logic_error("render domain prepared task is invalid");
|
|
if (stopping_.load(std::memory_order_acquire)) return Admission_Result::stopping;
|
|
{ std::lock_guard lock(task_mutex_); tasks_.push_back(std::move(task.task_)); }
|
|
task_condition_.notify_one();
|
|
const std::size_t queued = queued_.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_peak(peak_queued_, queued);
|
|
task.domain_.reset();
|
|
return Admission_Result::none;
|
|
}
|
|
Render_Domain::Statistics Render_Domain::statistics() const noexcept {
|
|
return {
|
|
static_cast<std::size_t>(default_capacity),
|
|
admitted_.load(std::memory_order_relaxed),
|
|
peak_admitted_.load(std::memory_order_relaxed),
|
|
queued_.load(std::memory_order_relaxed),
|
|
peak_queued_.load(std::memory_order_relaxed),
|
|
backpressure_count_.load(std::memory_order_relaxed),
|
|
backpressure_wait_ns_.load(std::memory_order_relaxed)
|
|
};
|
|
}
|
|
void Render_Domain::run() {
|
|
current_domain_ = this;
|
|
for (;;) {
|
|
std::unique_ptr<Task> task;
|
|
{ std::unique_lock lock(task_mutex_); task_condition_.wait(lock, [this] { return !tasks_.empty(); }); task = std::move(tasks_.front()); tasks_.pop_front(); }
|
|
if (!task) {
|
|
slots_.release();
|
|
current_domain_ = nullptr;
|
|
return;
|
|
}
|
|
queued_.fetch_sub(1, std::memory_order_relaxed);
|
|
release_admission();
|
|
try {
|
|
task->function();
|
|
}
|
|
catch (...) {
|
|
task->on_exception(contextual_exception("executing render domain task", std::current_exception()));
|
|
}
|
|
task.reset();
|
|
if (destroy_on_exit_.load(std::memory_order_acquire)) {
|
|
current_domain_ = nullptr;
|
|
if (thread_.joinable()) thread_.detach();
|
|
delete this;
|
|
return;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
|