321 lines
17 KiB
C++
321 lines
17 KiB
C++
#include "render_common.hpp"
|
|
#include <atomic>
|
|
#include <chrono>
|
|
#include <limits>
|
|
#include <mutex>
|
|
#include <stdexcept>
|
|
#include <taskflow/observer/interface.hpp>
|
|
namespace aethera {
|
|
namespace {
|
|
class Task_Observer : public tf::ObserverInterface {
|
|
private:
|
|
using Clock = std::chrono::steady_clock;
|
|
struct Task_Statistics {
|
|
std::atomic_size_t count{};
|
|
std::atomic_uint64_t total_time_ns{};
|
|
std::atomic_uint64_t min_time_ns{std::numeric_limits<std::uint64_t>::max()};
|
|
std::atomic_uint64_t max_time_ns{};
|
|
};
|
|
struct Worker_Statistics {
|
|
std::atomic_size_t task_count{};
|
|
std::atomic_size_t peak_queue_size{};
|
|
std::atomic_size_t max_queue_capacity{};
|
|
std::atomic_uint64_t task_time_ns{};
|
|
std::atomic_uint64_t busy_time_ns{};
|
|
std::atomic_uint64_t min_task_time_ns{std::numeric_limits<std::uint64_t>::max()};
|
|
std::atomic_uint64_t max_task_time_ns{};
|
|
};
|
|
std::vector<std::vector<Clock::time_point>> starts;
|
|
std::vector<Clock::time_point> worker_busy_starts;
|
|
std::unique_ptr<Worker_Statistics[]> worker_statistics;
|
|
std::size_t worker_statistics_count{};
|
|
std::array<Task_Statistics, tf::TASK_TYPES.size()> task_types;
|
|
std::atomic_size_t active_tasks{};
|
|
std::atomic_size_t peak_active_tasks{};
|
|
std::atomic_size_t active_workers{};
|
|
std::atomic_size_t peak_active_workers{};
|
|
std::atomic_size_t completed_tasks{};
|
|
std::atomic_size_t named_tasks{};
|
|
std::atomic_size_t peak_queue_size{};
|
|
std::atomic_size_t max_queue_capacity{};
|
|
std::atomic_size_t max_predecessors{};
|
|
std::atomic_size_t max_successors{};
|
|
std::atomic_size_t max_strong_dependencies{};
|
|
std::atomic_size_t max_weak_dependencies{};
|
|
std::atomic_uint64_t total_execution_time_ns{};
|
|
std::atomic_uint64_t worker_busy_time_ns{};
|
|
std::atomic_uint64_t first_task_time_ns{};
|
|
std::atomic_uint64_t last_task_time_ns{};
|
|
std::atomic_uint64_t longest_task_time_ns{};
|
|
std::atomic_size_t longest_task_hash{};
|
|
std::atomic<tf::TaskType> longest_task_type{tf::TaskType::UNDEFINED};
|
|
static std::uint64_t clock_ns(Clock::time_point value) noexcept {
|
|
return static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(value.time_since_epoch()).count());
|
|
}
|
|
template <typename T>
|
|
static void update_max(std::atomic<T>& value, T next) noexcept {
|
|
auto current = value.load(std::memory_order_relaxed);
|
|
while (current < next && !value.compare_exchange_weak(current, next, std::memory_order_relaxed)) {}
|
|
}
|
|
static void update_min(std::atomic_uint64_t& value, std::uint64_t next) noexcept {
|
|
auto current = value.load(std::memory_order_relaxed);
|
|
while (next < current && !value.compare_exchange_weak(current, next, std::memory_order_relaxed)) {}
|
|
}
|
|
static void update_first(std::atomic_uint64_t& value, std::uint64_t next) noexcept {
|
|
auto current = value.load(std::memory_order_relaxed);
|
|
while ((current == 0 || next < current) && !value.compare_exchange_weak(current, next, std::memory_order_relaxed)) {}
|
|
}
|
|
static std::size_t task_type_index(tf::TaskType type) noexcept {
|
|
return static_cast<std::size_t>(type);
|
|
}
|
|
public:
|
|
void set_up(std::size_t workers) override {
|
|
starts.resize(workers);
|
|
worker_busy_starts.resize(workers);
|
|
worker_statistics = std::make_unique<Worker_Statistics[]>(workers);
|
|
worker_statistics_count = workers;
|
|
for (auto& worker : starts) worker.reserve(8);
|
|
}
|
|
void on_entry(tf::WorkerView worker, tf::TaskView task) override {
|
|
auto now = Clock::now();
|
|
auto& worker_starts = starts[worker.id()];
|
|
if (worker_starts.empty()) {
|
|
worker_busy_starts[worker.id()] = now;
|
|
auto active = active_workers.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_max(peak_active_workers, active);
|
|
}
|
|
worker_starts.push_back(now);
|
|
auto active = active_tasks.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
update_max(peak_active_tasks, active);
|
|
update_max(peak_queue_size, worker.queue_size());
|
|
update_max(max_queue_capacity, worker.queue_capacity());
|
|
auto& worker_state = worker_statistics[worker.id()];
|
|
update_max(worker_state.peak_queue_size, worker.queue_size());
|
|
update_max(worker_state.max_queue_capacity, worker.queue_capacity());
|
|
update_max(max_predecessors, task.num_predecessors());
|
|
update_max(max_successors, task.num_successors());
|
|
update_max(max_strong_dependencies, task.num_strong_dependencies());
|
|
update_max(max_weak_dependencies, task.num_weak_dependencies());
|
|
if (!task.name().empty()) named_tasks.fetch_add(1, std::memory_order_relaxed);
|
|
update_first(first_task_time_ns, clock_ns(now));
|
|
}
|
|
void on_exit(tf::WorkerView worker, tf::TaskView task) override {
|
|
auto now = Clock::now();
|
|
auto& worker_starts = starts[worker.id()];
|
|
auto start = worker_starts.back();
|
|
worker_starts.pop_back();
|
|
auto elapsed = static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(now - start).count());
|
|
total_execution_time_ns.fetch_add(elapsed, std::memory_order_relaxed);
|
|
completed_tasks.fetch_add(1, std::memory_order_relaxed);
|
|
auto& worker_state = worker_statistics[worker.id()];
|
|
worker_state.task_count.fetch_add(1, std::memory_order_relaxed);
|
|
worker_state.task_time_ns.fetch_add(elapsed, std::memory_order_relaxed);
|
|
update_min(worker_state.min_task_time_ns, elapsed);
|
|
update_max(worker_state.max_task_time_ns, elapsed);
|
|
auto type_index = task_type_index(task.type());
|
|
if (type_index < task_types.size()) {
|
|
auto& type = task_types[type_index];
|
|
type.count.fetch_add(1, std::memory_order_relaxed);
|
|
type.total_time_ns.fetch_add(elapsed, std::memory_order_relaxed);
|
|
update_min(type.min_time_ns, elapsed);
|
|
update_max(type.max_time_ns, elapsed);
|
|
}
|
|
auto longest = longest_task_time_ns.load(std::memory_order_relaxed);
|
|
if (longest < elapsed && longest_task_time_ns.compare_exchange_strong(
|
|
longest, elapsed, std::memory_order_relaxed)) {
|
|
longest_task_hash.store(task.hash_value(), std::memory_order_relaxed);
|
|
longest_task_type.store(task.type(), std::memory_order_relaxed);
|
|
}
|
|
active_tasks.fetch_sub(1, std::memory_order_relaxed);
|
|
if (worker_starts.empty()) {
|
|
auto busy = static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(now - worker_busy_starts[worker.id()]).count());
|
|
worker_busy_time_ns.fetch_add(busy, std::memory_order_relaxed);
|
|
worker_state.busy_time_ns.fetch_add(busy, std::memory_order_relaxed);
|
|
active_workers.fetch_sub(1, std::memory_order_relaxed);
|
|
}
|
|
update_max(last_task_time_ns, clock_ns(now));
|
|
}
|
|
void write_state(Task_Runtime_State& state, std::size_t workers, std::size_t active_topologies) const {
|
|
state.worker_count = workers;
|
|
state.active_topology_count = active_topologies;
|
|
state.active_task_count = active_tasks.load(std::memory_order_relaxed);
|
|
state.peak_active_task_count = peak_active_tasks.load(std::memory_order_relaxed);
|
|
state.active_worker_count = active_workers.load(std::memory_order_relaxed);
|
|
state.peak_active_worker_count = peak_active_workers.load(std::memory_order_relaxed);
|
|
state.observed_task_count = completed_tasks.load(std::memory_order_relaxed);
|
|
state.named_task_count = named_tasks.load(std::memory_order_relaxed);
|
|
state.peak_observed_worker_queue_size = peak_queue_size.load(std::memory_order_relaxed);
|
|
state.max_observed_worker_queue_capacity = max_queue_capacity.load(std::memory_order_relaxed);
|
|
state.max_predecessors = max_predecessors.load(std::memory_order_relaxed);
|
|
state.max_successors = max_successors.load(std::memory_order_relaxed);
|
|
state.max_strong_dependencies = max_strong_dependencies.load(std::memory_order_relaxed);
|
|
state.max_weak_dependencies = max_weak_dependencies.load(std::memory_order_relaxed);
|
|
state.total_task_time_ns = total_execution_time_ns.load(std::memory_order_relaxed);
|
|
state.worker_busy_time_ns = worker_busy_time_ns.load(std::memory_order_relaxed);
|
|
auto first = first_task_time_ns.load(std::memory_order_relaxed);
|
|
auto last = last_task_time_ns.load(std::memory_order_relaxed);
|
|
state.observed_wall_time_ns = first && last >= first ? last - first : 0;
|
|
state.worker_utilization = workers && state.observed_wall_time_ns ? static_cast<double>(state.worker_busy_time_ns) * 100.0 / static_cast<double>(state.observed_wall_time_ns) / static_cast<double>(workers) : 0.0;
|
|
for (std::size_t i = 0; i < task_types.size(); ++i) {
|
|
const auto& source = task_types[i];
|
|
auto& target = state.task_types[i];
|
|
target.count = source.count.load(std::memory_order_relaxed);
|
|
target.total_time_ns = source.total_time_ns.load(std::memory_order_relaxed);
|
|
auto min = source.min_time_ns.load(std::memory_order_relaxed);
|
|
target.min_time_ns = target.count ? min : 0;
|
|
target.max_time_ns = source.max_time_ns.load(std::memory_order_relaxed);
|
|
}
|
|
state.workers.resize(worker_statistics_count);
|
|
for (std::size_t i = 0; i < worker_statistics_count; ++i) {
|
|
const auto& source = worker_statistics[i];
|
|
auto& target = state.workers[i];
|
|
target.id = i;
|
|
target.task_count = source.task_count.load(std::memory_order_relaxed);
|
|
target.peak_observed_queue_size = source.peak_queue_size.load(std::memory_order_relaxed);
|
|
target.max_observed_queue_capacity = source.max_queue_capacity.load(std::memory_order_relaxed);
|
|
target.task_time_ns = source.task_time_ns.load(std::memory_order_relaxed);
|
|
target.busy_time_ns = source.busy_time_ns.load(std::memory_order_relaxed);
|
|
target.idle_time_ns = state.observed_wall_time_ns > target.busy_time_ns ? state.observed_wall_time_ns - target.busy_time_ns : 0;
|
|
auto min = source.min_task_time_ns.load(std::memory_order_relaxed);
|
|
target.min_task_time_ns = target.task_count ? min : 0;
|
|
target.max_task_time_ns = source.max_task_time_ns.load(std::memory_order_relaxed);
|
|
target.utilization = state.observed_wall_time_ns ? static_cast<double>(target.busy_time_ns) * 100.0 / static_cast<double>(state.observed_wall_time_ns) : 0.0;
|
|
}
|
|
state.longest_task_time_ns = longest_task_time_ns.load(std::memory_order_relaxed);
|
|
state.longest_task_hash = longest_task_hash.load(std::memory_order_relaxed);
|
|
state.longest_task_name.clear();
|
|
state.longest_task_type = longest_task_type.load(std::memory_order_relaxed);
|
|
}
|
|
};
|
|
class Task_Resource : Pinned {
|
|
private:
|
|
std::unique_ptr<tf::Executor> executor;
|
|
std::shared_ptr<Task_Observer> observer;
|
|
std::pmr::memory_resource* memory;
|
|
std::atomic_size_t active_taskflows{};
|
|
std::atomic_size_t peak_active_taskflows{};
|
|
std::atomic_size_t completed_taskflows{};
|
|
std::atomic_size_t failed_taskflows{};
|
|
Task_Runtime_State state;
|
|
double_buffer::detail::State_Callback_Storage<Task_Runtime_State> state_callbacks;
|
|
std::recursive_mutex state_mutex;
|
|
std::atomic_bool state_callback_enabled{};
|
|
void create_executor(std::size_t workers, std::shared_ptr<tf::WorkerInterface> worker_interface) {
|
|
executor = std::make_unique<tf::Executor>(workers, std::move(worker_interface));
|
|
observer = executor->make_observer<Task_Observer>();
|
|
state = {};
|
|
}
|
|
void ensure_executor() {
|
|
std::lock_guard guard(state_mutex);
|
|
if (executor) return;
|
|
executor = std::make_unique<tf::Executor>();
|
|
observer = executor->make_observer<Task_Observer>();
|
|
}
|
|
void publish_state() {
|
|
if (!state_callback_enabled.load(std::memory_order_acquire)) return;
|
|
std::lock_guard guard(state_mutex);
|
|
observer->write_state(state, executor->num_workers(), executor->num_topologies());
|
|
state.active_taskflow_count = active_taskflows.load(std::memory_order_relaxed);
|
|
state.peak_active_taskflow_count = peak_active_taskflows.load(std::memory_order_relaxed);
|
|
state.completed_taskflow_count = completed_taskflows.load(std::memory_order_relaxed);
|
|
state.failed_taskflow_count = failed_taskflows.load(std::memory_order_relaxed);
|
|
state_callbacks.template notify<Task_Runtime_State_Tag>(state);
|
|
}
|
|
public:
|
|
Task_Resource() : memory(std::pmr::get_default_resource()) {}
|
|
static Task_Resource& instance() {
|
|
static Task_Resource value;
|
|
return value;
|
|
}
|
|
void initialize(std::size_t workers, std::shared_ptr<tf::WorkerInterface> worker_interface, Pmr pmr) {
|
|
std::lock_guard guard(state_mutex);
|
|
memory = pmr.resource();
|
|
create_executor(workers, std::move(worker_interface));
|
|
active_taskflows.store(0, std::memory_order_relaxed);
|
|
peak_active_taskflows.store(0, std::memory_order_relaxed);
|
|
completed_taskflows.store(0, std::memory_order_relaxed);
|
|
failed_taskflows.store(0, std::memory_order_relaxed);
|
|
}
|
|
std::uint64_t run(tf::Taskflow& taskflow) {
|
|
ensure_executor();
|
|
auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
auto peak = peak_active_taskflows.load(std::memory_order_relaxed);
|
|
while (peak < active && !peak_active_taskflows.compare_exchange_weak(peak, active, std::memory_order_relaxed)) {}
|
|
auto start = std::chrono::steady_clock::now();
|
|
try {
|
|
if (executor->this_worker())
|
|
executor->corun(taskflow);
|
|
else
|
|
executor->run(taskflow).get();
|
|
}
|
|
catch (...) {
|
|
active_taskflows.fetch_sub(1, std::memory_order_relaxed);
|
|
failed_taskflows.fetch_add(1, std::memory_order_relaxed);
|
|
publish_state();
|
|
throw;
|
|
}
|
|
auto elapsed = static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(std::chrono::steady_clock::now() - start).count());
|
|
active_taskflows.fetch_sub(1, std::memory_order_relaxed);
|
|
completed_taskflows.fetch_add(1, std::memory_order_relaxed);
|
|
publish_state();
|
|
return elapsed;
|
|
}
|
|
void run(tf::Taskflow& taskflow, std::function<void()> completion) {
|
|
if (!completion) throw std::invalid_argument("Taskflow completion is empty");
|
|
ensure_executor();
|
|
auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1;
|
|
auto peak = peak_active_taskflows.load(std::memory_order_relaxed);
|
|
while (peak < active && !peak_active_taskflows.compare_exchange_weak(
|
|
peak, active, std::memory_order_relaxed)) {}
|
|
executor->run(taskflow, [this, completion = std::move(completion)]() mutable {
|
|
active_taskflows.fetch_sub(1, std::memory_order_relaxed);
|
|
completed_taskflows.fetch_add(1, std::memory_order_relaxed);
|
|
publish_state();
|
|
completion();
|
|
});
|
|
}
|
|
void schedule(std::function<void()> task) {
|
|
if (!task) throw std::invalid_argument("Taskflow scheduled task is empty");
|
|
ensure_executor();
|
|
executor->silent_async(std::move(task));
|
|
}
|
|
std::pmr::memory_resource* memory_resource() const noexcept {
|
|
return memory;
|
|
}
|
|
void set_state_callback(std::function<void(const Task_Runtime_State&)> callback) {
|
|
std::lock_guard guard(state_mutex);
|
|
state_callbacks.template set<Task_Runtime_State_Tag>(std::move(callback));
|
|
state_callback_enabled.store(true, std::memory_order_release);
|
|
}
|
|
void clear_state_callback() {
|
|
std::lock_guard guard(state_mutex);
|
|
state_callbacks.template clear<Task_Runtime_State_Tag>();
|
|
state_callback_enabled.store(false, std::memory_order_release);
|
|
}
|
|
};
|
|
}
|
|
void initialize_runtime(std::size_t workers, std::shared_ptr<tf::WorkerInterface> worker_interface, Pmr pmr) {
|
|
Task_Resource::instance().initialize(workers, std::move(worker_interface), pmr);
|
|
}
|
|
void schedule_task(std::function<void()> task) {
|
|
Task_Resource::instance().schedule(std::move(task));
|
|
}
|
|
namespace detail {
|
|
void set_runtime_state_callback_impl(std::function<void(const Task_Runtime_State&)> callback) {
|
|
Task_Resource::instance().set_state_callback(std::move(callback));
|
|
}
|
|
void clear_runtime_state_callback_impl() {
|
|
Task_Resource::instance().clear_state_callback();
|
|
}
|
|
std::uint64_t run_taskflow(tf::Taskflow& taskflow) {
|
|
return Task_Resource::instance().run(taskflow);
|
|
}
|
|
void run_taskflow(tf::Taskflow& taskflow, std::function<void()> completion) {
|
|
Task_Resource::instance().run(taskflow, std::move(completion));
|
|
}
|
|
std::pmr::memory_resource* task_memory_resource() noexcept {
|
|
return Task_Resource::instance().memory_resource();
|
|
}
|
|
}
|
|
}
|