优化了 还是卡

This commit is contained in:
2026-08-25 15:03:33 +08:00
parent 2c3105cc07
commit fa93927d81
23 changed files with 1517 additions and 619 deletions
+48 -2
View File
@@ -1,5 +1,7 @@
#include "Task_Graph.hpp"
#include "Task_Graph_Internal.hpp"
#include <algorithm>
#include <cstdio>
#include <stdexcept>
#include <taskflow/taskflow.hpp>
#include <unordered_map>
@@ -12,6 +14,7 @@ struct Task_Graph::Private {
std::string name{}; /* 调用方提供的业务节点名。 */
std::string node_id{}; /* 本图内去重后的业务节点 ID。 */
std::shared_ptr<Private> child{}; /* 模块引用的子图;普通节点为空。 */
std::vector<std::pair<std::string, std::string>> attributes{};
};
explicit Private(std::string graph_name)
: taskflow(std::move(graph_name)) {}
@@ -43,6 +46,7 @@ struct Task_Graph::Private {
node.parent_node_id = std::string(parent);
node.name = source.name;
node.type = std::string(tf::to_string(source.task.type()));
node.attributes = source.attributes;
source.task.for_each_predecessor([&](tf::Task value) {
node.predecessors.push_back(
static_cast<std::uint64_t>(value.hash_value()));
@@ -98,6 +102,22 @@ void Task_Node::precede(const Task_Node& after) const {
d->graph->nodes[after.d->index].task);
}
Task_Node& Task_Node::describe(std::string key, std::string value) {
if (!d || !d->graph || d->generation != d->graph->generation ||
d->index >= d->graph->nodes.size())
throw std::invalid_argument("Task graph node handle is stale");
if (key.empty())
throw std::invalid_argument("Task graph node attribute key is empty");
auto& attributes = d->graph->nodes[d->index].attributes;
const auto found = std::ranges::find(
attributes, key, &std::pair<std::string, std::string>::first);
if (found == attributes.end())
attributes.emplace_back(std::move(key), std::move(value));
else
found->second = std::move(value);
return *this;
}
Task_Graph::Task_Graph(std::string graph_name)
: d(std::make_shared<Private>(std::move(graph_name))) {}
Task_Graph::~Task_Graph() = default;
@@ -122,7 +142,20 @@ Task_Graph& Task_Graph::operator=(Task_Graph&& other) noexcept {
Task_Node Task_Graph::add(std::string task_name,
std::function<void()> work) {
if (!work) throw std::invalid_argument("Task graph work is empty");
auto task = d->taskflow.emplace(std::move(work));
auto task = d->taskflow.emplace(
[name = task_name, work = std::move(work)] {
try { work(); }
catch (const std::exception& error) {
std::fprintf(stderr, "Taskflow node %s failed: %s\n",
name.c_str(), error.what());
throw;
}
catch (...) {
std::fprintf(stderr, "Taskflow node %s failed: non-standard exception\n",
name.c_str());
throw;
}
});
const auto node_id = d->unique_node_id(task_name);
task.name(node_id);
d->nodes.push_back({task, std::move(task_name), node_id, nullptr});
@@ -133,7 +166,20 @@ Task_Node Task_Graph::add(std::string task_name,
Task_Node Task_Graph::add_condition(std::string task_name,
std::function<int()> work) {
if (!work) throw std::invalid_argument("Task graph condition is empty");
auto task = d->taskflow.emplace(std::move(work));
auto task = d->taskflow.emplace(
[name = task_name, work = std::move(work)] {
try { return work(); }
catch (const std::exception& error) {
std::fprintf(stderr, "Taskflow condition %s failed: %s\n",
name.c_str(), error.what());
throw;
}
catch (...) {
std::fprintf(stderr, "Taskflow condition %s failed: non-standard exception\n",
name.c_str());
throw;
}
});
const auto node_id = d->unique_node_id(task_name);
task.name(node_id);
d->nodes.push_back({task, std::move(task_name), node_id, nullptr});
+1
View File
@@ -19,6 +19,7 @@ public:
Task_Node& operator=(Task_Node&&) noexcept;
/* 建立本节点到 after 的有向依赖;两个节点必须属于同一张图。 */
void precede(const Task_Node& after) const;
Task_Node& describe(std::string key, std::string value);
private:
friend class Task_Graph;
struct Private;
+10 -4
View File
@@ -18,10 +18,16 @@ struct Taskflow_Frame_Access {
static void release_writer(Render_Frame& frame) noexcept;
[[nodiscard]] static bool contains_task(const Render_Frame& frame,
std::uint64_t native_id) noexcept;
static void append_task(Render_Frame& frame, std::size_t worker,
std::uint64_t native_id, std::size_t queue_size,
std::size_t queue_capacity, Clock::time_point started,
Clock::time_point finished);
[[nodiscard]] static std::size_t append_task(
Render_Frame& frame, std::size_t worker, std::uint64_t native_id,
std::size_t queue_size, std::size_t queue_capacity,
Clock::time_point entered, Clock::time_point started,
Clock::time_point finished, std::uint64_t cpu_entered_ns,
std::uint64_t cpu_started_ns, std::uint64_t cpu_finished_ns);
static void finish_task_observer(
Render_Frame& frame, std::size_t worker, std::size_t task,
Clock::time_point completed, std::uint64_t cpu_finished_ns,
std::uint64_t cpu_completed_ns) noexcept;
[[nodiscard]] static Taskflow_Graph_Token begin_graph(
Render_Frame& frame, Task_Graph& graph, std::string_view stage);
static void finish_graph(Taskflow_Graph_Token token) noexcept;
+135 -20
View File
@@ -6,12 +6,14 @@
#include <array>
#include <atomic>
#include <chrono>
#include <functional>
#include <limits>
#include <mutex>
#include <optional>
#include <ranges>
#include <thread>
#include <unordered_map>
#include <unordered_set>
namespace aethera {
namespace {
constexpr std::size_t marker_count = static_cast<std::size_t>(Frame_Trace_Marker::count);
@@ -174,7 +176,28 @@ Frame_Statistics_Sample Render_Frame::statistics(Frame_Dimension dimension) cons
if (dimension == Frame_Dimension::two_dimensional) {
result.set(Frame_Statistic::pipeline_2d_event_ms, take(event_time));
result.set(Frame_Statistic::pipeline_2d_prepare_ms, take(prepare_time));
result.set(Frame_Statistic::pipeline_2d_paint_ms, take(paint_time));
result.set(Frame_Statistic::pipeline_2d_paint_ms, paint_time);
double paint_remaining = paint_time;
const auto take_paint = [&](Frame_Trace_Marker first,
Frame_Trace_Marker last) {
const double value = std::min(paint_remaining, interval(first, last));
paint_remaining -= value;
return take(value);
};
result.set(Frame_Statistic::pipeline_2d_frame_target_ms, take_paint(
Frame_Trace_Marker::paint_frame_target_started,
Frame_Trace_Marker::paint_frame_target_finished));
result.set(Frame_Statistic::pipeline_2d_background_ms, take_paint(
Frame_Trace_Marker::paint_background_started,
Frame_Trace_Marker::paint_background_finished));
result.set(Frame_Statistic::pipeline_2d_cache_targets_ms, take_paint(
Frame_Trace_Marker::paint_cache_targets_started,
Frame_Trace_Marker::paint_cache_targets_finished));
result.set(Frame_Statistic::pipeline_2d_taskflow_ms, take_paint(
Frame_Trace_Marker::paint_taskflow_started,
Frame_Trace_Marker::paint_taskflow_finished));
result.set(Frame_Statistic::pipeline_2d_paint_coordination_ms,
take(paint_remaining));
result.set(Frame_Statistic::pipeline_2d_scene_coordination_ms,
take(std::max(0.0, scene_time - event_time - prepare_time - paint_time)));
result.set(Frame_Statistic::pipeline_2d_callback_ms, take(interval(
@@ -255,6 +278,12 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const {
result.identity = d->identity;
result.created_time_unix_ns = d->created_time_unix_ns;
result.worker_count = d->taskflow_workers.size();
for (std::size_t index = 0; index < marker_count; ++index) {
const auto encoded = d->markers[index].load(std::memory_order_acquire);
if (!encoded) continue;
result.markers.push_back(Frame_Trace_Point{
static_cast<Frame_Trace_Marker>(index), decode_present_value(encoded)});
}
{
std::lock_guard lock(d->taskflow_graph_mutex);
result.graphs = d->taskflow_graphs;
@@ -271,24 +300,80 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const {
});
});
std::unordered_map<std::uint64_t, double> finished;
struct Node_Context {
const Taskflow_Graph_Trace* graph{};
const Taskflow_Graph_Trace::Node* node{};
};
std::unordered_map<std::uint64_t, Node_Context> nodes;
std::unordered_map<std::string, std::uint64_t> native_id_by_node_id;
std::unordered_map<std::string, std::vector<std::uint64_t>> children_by_parent_id;
for (const auto& graph : result.graphs) {
for (const auto& node : graph.nodes) {
nodes.try_emplace(node.native_id, Node_Context{&graph, &node});
native_id_by_node_id.try_emplace(node.node_id, node.native_id);
if (!node.parent_node_id.empty())
children_by_parent_id[node.parent_node_id].push_back(node.native_id);
}
}
std::unordered_map<std::uint64_t, double> observed_finished;
for (const auto& task : result.tasks)
observed_finished[task.native_id] = std::max(
observed_finished[task.native_id], task.completed_ms);
std::unordered_map<std::uint64_t, double> ready_cache;
std::unordered_map<std::uint64_t, double> finished_cache;
std::unordered_set<std::uint64_t> resolving_ready;
std::unordered_set<std::uint64_t> resolving_finished;
std::function<double(std::uint64_t)> ready_of;
std::function<double(std::uint64_t)> finished_of;
ready_of = [&](std::uint64_t native_id) -> double {
if (const auto cached = ready_cache.find(native_id);
cached != ready_cache.end()) return cached->second;
const auto context = nodes.find(native_id);
if (context == nodes.end() || !resolving_ready.insert(native_id).second)
return 0.0;
double ready = context->second.graph->submitted_ms;
const auto& node = *context->second.node;
if (!node.parent_node_id.empty()) {
if (const auto parent = native_id_by_node_id.find(node.parent_node_id);
parent != native_id_by_node_id.end())
ready = std::max(ready, ready_of(parent->second));
}
for (const auto predecessor : node.predecessors)
ready = std::max(ready, finished_of(predecessor));
resolving_ready.erase(native_id);
ready_cache.emplace(native_id, ready);
return ready;
};
finished_of = [&](std::uint64_t native_id) -> double {
if (const auto cached = finished_cache.find(native_id);
cached != finished_cache.end()) return cached->second;
if (!resolving_finished.insert(native_id).second) return ready_of(native_id);
double finished = ready_of(native_id);
if (const auto observed = observed_finished.find(native_id);
observed != observed_finished.end())
finished = std::max(finished, observed->second);
if (const auto context = nodes.find(native_id); context != nodes.end()) {
if (const auto children = children_by_parent_id.find(
context->second.node->node_id);
children != children_by_parent_id.end())
for (const auto child : children->second)
finished = std::max(finished, finished_of(child));
}
resolving_finished.erase(native_id);
finished_cache.emplace(native_id, finished);
return finished;
};
std::unordered_map<std::uint64_t, double> previous_execution_finished;
for (auto& task : result.tasks) {
double ready{};
bool matched_graph{};
for (const auto& graph : result.graphs)
for (const auto& node : graph.nodes)
if (node.native_id == task.native_id) {
if (!matched_graph || graph.submitted_ms < ready)
ready = graph.submitted_ms;
matched_graph = true;
for (const auto predecessor : node.predecessors)
if (const auto found = finished.find(predecessor);
found != finished.end())
ready = std::max(ready, found->second);
}
double ready = ready_of(task.native_id);
if (const auto previous = previous_execution_finished.find(task.native_id);
previous != previous_execution_finished.end())
ready = std::max(ready, previous->second);
task.ready_ms = ready;
task.queue_wait_ms = std::max(0.0, task.started_ms - ready);
finished[task.native_id] = std::max(finished[task.native_id], task.finished_ms);
task.queue_wait_ms = std::max(0.0, task.entered_ms - ready);
previous_execution_finished[task.native_id] = task.completed_ms;
}
return result;
}
@@ -344,12 +429,15 @@ bool detail::Taskflow_Frame_Access::contains_task(
}
}
void detail::Taskflow_Frame_Access::append_task(
std::size_t detail::Taskflow_Frame_Access::append_task(
Render_Frame& frame, std::size_t worker, std::uint64_t native_id,
std::size_t queue_size, std::size_t queue_capacity,
Clock::time_point started, Clock::time_point finished) {
Clock::time_point entered, Clock::time_point started,
Clock::time_point finished, std::uint64_t cpu_entered_ns,
std::uint64_t cpu_started_ns, std::uint64_t cpu_finished_ns) {
auto& data = *frame.d;
if (worker >= data.taskflow_workers.size()) return;
if (worker >= data.taskflow_workers.size())
return std::numeric_limits<std::size_t>::max();
const auto elapsed_ms = [&](Clock::time_point value) {
return std::chrono::duration<double, std::milli>(value - data.created_at).count();
};
@@ -358,10 +446,37 @@ void detail::Taskflow_Frame_Access::append_task(
trace.worker_id = worker;
trace.worker_queue_size = queue_size;
trace.worker_queue_capacity = queue_capacity;
trace.entered_ms = elapsed_ms(entered);
trace.started_ms = elapsed_ms(started);
trace.finished_ms = elapsed_ms(finished);
trace.completed_ms = trace.finished_ms;
trace.duration_ms = std::max(0.0, trace.finished_ms - trace.started_ms);
trace.cpu_duration_ms = cpu_finished_ns >= cpu_started_ns
? static_cast<double>(cpu_finished_ns - cpu_started_ns) / 1'000'000.0
: 0.0;
trace.observer_entry_ms = std::max(0.0, trace.started_ms - trace.entered_ms);
trace.observer_entry_cpu_ms = cpu_started_ns >= cpu_entered_ns
? static_cast<double>(cpu_started_ns - cpu_entered_ns) / 1'000'000.0
: 0.0;
data.taskflow_workers[worker].tasks.push_back(std::move(trace));
return data.taskflow_workers[worker].tasks.size() - 1;
}
void detail::Taskflow_Frame_Access::finish_task_observer(
Render_Frame& frame, std::size_t worker, std::size_t task,
Clock::time_point completed, std::uint64_t cpu_finished_ns,
std::uint64_t cpu_completed_ns) noexcept {
auto& data = *frame.d;
if (worker >= data.taskflow_workers.size() ||
task >= data.taskflow_workers[worker].tasks.size()) return;
auto& trace = data.taskflow_workers[worker].tasks[task];
trace.completed_ms = std::chrono::duration<double, std::milli>(
completed - data.created_at).count();
trace.observer_exit_ms = std::max(
0.0, trace.completed_ms - trace.finished_ms);
trace.observer_exit_cpu_ms = cpu_completed_ns >= cpu_finished_ns
? static_cast<double>(cpu_completed_ns - cpu_finished_ns) / 1'000'000.0
: 0.0;
}
detail::Taskflow_Graph_Token detail::Taskflow_Frame_Access::begin_graph(
+21 -4
View File
@@ -17,6 +17,14 @@ enum class Frame_Trace_Marker : std::uint8_t {
prepare_started,
prepare_finished,
paint_started,
paint_frame_target_started,
paint_frame_target_finished,
paint_background_started,
paint_background_finished,
paint_cache_targets_started,
paint_cache_targets_finished,
paint_taskflow_started,
paint_taskflow_finished,
paint_finished,
backend_queue_entered,
backend_prepare_started,
@@ -73,6 +81,7 @@ struct Taskflow_Graph_Trace {
std::string type{}; /* Taskflow 原生 TaskType。 */
std::vector<std::uint64_t> predecessors{}; /* 原生直接前驱 hash。 */
std::vector<std::uint64_t> successors{}; /* 原生直接后继 hash。 */
std::vector<std::pair<std::string, std::string>> attributes{};
};
std::vector<Node> nodes{}; /* DAG 构造时登记的节点与依赖元信息。 */
double submitted_ms{}; /* 相对帧创建时刻的 run 提交时间。 */
@@ -84,16 +93,24 @@ struct Taskflow_Task_Trace {
std::size_t worker_id{}; /* 执行该任务的 Executor worker。 */
std::size_t worker_queue_size{}; /* on_entry 时的原生 worker queue_size。 */
std::size_t worker_queue_capacity{}; /* on_entry 时的原生 worker queue_capacity。 */
double started_ms{}; /* 相对帧创建时刻的 on_entry 时间。 */
double finished_ms{}; /* 相对帧创建时刻的 on_exit 时间。 */
double duration_ms{}; /* 原生 on_entry 到 on_exit 的持续时间。 */
double entered_ms{}; /* Observer on_entry 开始,即 Executor 真正选中任务的时间。 */
double started_ms{}; /* Observer on_entry 返回前,即任务体即将执行的时间。 */
double finished_ms{}; /* Observer on_exit 进入,即任务体已经结束的时间。 */
double completed_ms{}; /* Observer on_exit 与按帧追踪写入全部结束的时间。 */
double duration_ms{}; /* 仅任务体 started 到 finished 的持续时间。 */
double cpu_duration_ms{}; /* 任务体在当前 worker 线程上实际消耗的 CPU 时间。 */
double observer_entry_ms{}; /* on_entry 诊断本身的耗时。 */
double observer_exit_ms{}; /* on_exit 诊断与按帧追踪写入的耗时。 */
double observer_entry_cpu_ms{}; /* on_entry 诊断实际消耗的 worker CPU 时间。 */
double observer_exit_cpu_ms{}; /* on_exit 诊断实际消耗的 worker CPU 时间。 */
double ready_ms{}; /* 前驱完成或根 run 提交后的估算就绪时间。 */
double queue_wait_ms{}; /* ready 到 on_entry 的估算 Executor 排队时间。 */
double queue_wait_ms{}; /* ready 到 entered 的估算 Executor 排队时间。 */
};
struct Taskflow_Frame_Trace {
Frame_Identity identity{}; /* 该执行图所属逻辑渲染帧。 */
std::uint64_t created_time_unix_ns{}; /* 帧创建 Unix 时间,单位纳秒。 */
std::size_t worker_count{}; /* 捕获时全局 Executor 的 worker 数。 */
std::vector<Frame_Trace_Point> markers{}; /* 与该帧 DAG 共用时间原点的原始流水线时间点。 */
std::vector<Taskflow_Graph_Trace> graphs{}; /* 本帧主动执行的业务 DAG 元信息。 */
std::vector<Taskflow_Task_Trace> tasks{}; /* 本帧窗口内原生 Observer 完成的任务执行。 */
};
+5
View File
@@ -36,6 +36,11 @@ enum class Frame_Statistic : std::uint8_t {
pipeline_2d_event_ms,
pipeline_2d_prepare_ms,
pipeline_2d_paint_ms,
pipeline_2d_frame_target_ms,
pipeline_2d_background_ms,
pipeline_2d_cache_targets_ms,
pipeline_2d_taskflow_ms,
pipeline_2d_paint_coordination_ms,
pipeline_2d_scene_coordination_ms,
pipeline_2d_callback_ms,
pipeline_2d_frame_handoff_ms,
+72 -20
View File
@@ -5,10 +5,18 @@
#include <chrono>
#include <limits>
#include <mutex>
#include <optional>
#include <shared_mutex>
#include <stdexcept>
#include <taskflow/observer/interface.hpp>
#include <taskflow/taskflow.hpp>
#if defined(_WIN32)
#define WIN32_LEAN_AND_MEAN
#define NOMINMAX
#include <Windows.h>
#elif defined(CLOCK_THREAD_CPUTIME_ID)
#include <ctime>
#endif
namespace aethera {
namespace {
tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept {
@@ -16,6 +24,28 @@ tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept {
detail::Task_Graph_Access::native_storage(graph));
}
std::uint64_t current_thread_cpu_ns() noexcept {
#if defined(_WIN32)
FILETIME created{}, exited{}, kernel{}, user{};
if (!GetThreadTimes(GetCurrentThread(), &created, &exited, &kernel, &user))
return 0;
const auto ticks = [](FILETIME value) noexcept {
ULARGE_INTEGER result{};
result.LowPart = value.dwLowDateTime;
result.HighPart = value.dwHighDateTime;
return result.QuadPart;
};
return (ticks(kernel) + ticks(user)) * 100ULL;
#elif defined(CLOCK_THREAD_CPUTIME_ID)
timespec value{};
if (clock_gettime(CLOCK_THREAD_CPUTIME_ID, &value) != 0) return 0;
return static_cast<std::uint64_t>(value.tv_sec) * 1'000'000'000ULL +
static_cast<std::uint64_t>(value.tv_nsec);
#else
return 0;
#endif
}
class Task_Observer : public tf::ObserverInterface {
private:
using Clock = std::chrono::steady_clock;
@@ -40,7 +70,10 @@ private:
std::atomic_uint64_t max_task_time_ns{};
};
struct Start_Record {
Clock::time_point started{}; /* Observer on_entry 时间。 */
Clock::time_point entered{}; /* Observer on_entry 进入时间。 */
Clock::time_point started{}; /* on_entry 完成、任务体即将执行的时间。 */
std::uint64_t cpu_entered_ns{}; /* on_entry 进入时的 worker CPU 时间。 */
std::uint64_t cpu_started_ns{}; /* 任务体开始前的 worker CPU 时间。 */
Render_Frame* frame{}; /* 进入任务时唯一活动的按帧捕获。 */
std::size_t queue_size{}; /* 进入任务时 worker 队列深度。 */
std::size_t queue_capacity{}; /* 进入任务时 worker 队列容量。 */
@@ -110,7 +143,8 @@ public:
update_max(peak_active_workers, active);
}
worker_starts.push_back(Start_Record{
now, nullptr, worker.queue_size(), worker.queue_capacity(),
now, {}, current_thread_cpu_ns(), 0, nullptr,
worker.queue_size(), worker.queue_capacity(),
static_cast<std::uint64_t>(task.hash_value()), task.type()});
auto* frame = trace_frame.load(std::memory_order_acquire);
/*
@@ -148,13 +182,18 @@ public:
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));
worker_starts.back().cpu_started_ns = current_thread_cpu_ns();
worker_starts.back().started = Clock::now();
}
void on_exit(tf::WorkerView worker, tf::TaskView task) override {
auto now = Clock::now();
const auto finished = Clock::now();
const auto cpu_finished_ns = current_thread_cpu_ns();
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.started).count());
const auto elapsed = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::nanoseconds>(
finished - start.started).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()];
@@ -180,8 +219,33 @@ public:
std::memory_order_release);
}
active_tasks.fetch_sub(1, std::memory_order_relaxed);
std::optional<std::size_t> trace_task;
if (start.frame) {
try {
trace_task = detail::Taskflow_Frame_Access::append_task(
*start.frame, worker.id(),
static_cast<std::uint64_t>(task.hash_value()),
start.queue_size, start.queue_capacity, start.entered,
start.started, finished, start.cpu_entered_ns,
start.cpu_started_ns, cpu_finished_ns);
}
catch (...) {
/* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */
}
}
const auto cpu_completed_ns = current_thread_cpu_ns();
const auto completed = Clock::now();
if (start.frame) {
if (trace_task)
detail::Taskflow_Frame_Access::finish_task_observer(
*start.frame, worker.id(), *trace_task, completed,
cpu_finished_ns, cpu_completed_ns);
detail::Taskflow_Frame_Access::release_writer(*start.frame);
}
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());
auto busy = static_cast<std::uint64_t>(
std::chrono::duration_cast<std::chrono::nanoseconds>(
completed - 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);
@@ -195,27 +259,15 @@ public:
else {
const auto& parent = worker_starts.back();
worker_state.active_task_hash.store(parent.native_id, std::memory_order_relaxed);
worker_state.active_task_started_ns.store(clock_ns(parent.started),
worker_state.active_task_started_ns.store(clock_ns(parent.entered),
std::memory_order_relaxed);
worker_state.active_task_type.store(parent.type, std::memory_order_relaxed);
worker_state.current_queue_size.store(parent.queue_size,
std::memory_order_relaxed);
worker_state.current_queue_capacity.store(parent.queue_capacity,
std::memory_order_relaxed);
}
update_max(last_task_time_ns, clock_ns(now));
if (start.frame) {
try {
detail::Taskflow_Frame_Access::append_task(
*start.frame, worker.id(),
static_cast<std::uint64_t>(task.hash_value()),
start.queue_size, start.queue_capacity, start.started, now);
}
catch (...) {
/* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */
}
detail::Taskflow_Frame_Access::release_writer(*start.frame);
std::memory_order_relaxed);
}
update_max(last_task_time_ns, clock_ns(completed));
}
bool begin_trace(Render_Frame& frame, std::size_t workers) {
std::unique_lock guard(trace_mutex);
+9 -4
View File
@@ -164,8 +164,6 @@ void Scene::Private::after_advance(Object* object,
if (dispatch->prepare.run) dispatch->prepare.run(root);
});
}
auto prepare_extension = taskflow.compose(task_prefix + ".extension",
data->prepare_extension);
auto prepare_done = taskflow.add(task_prefix + ".complete", [dispatch, root] {
auto& state = *dispatch->state.pending(root);
if (state.prepare_executed) {
@@ -179,8 +177,15 @@ void Scene::Private::after_advance(Object* object,
});
prepare_if.precede(prepare_run);
prepare_if.precede(prepare_done);
prepare_run.precede(prepare_extension);
prepare_extension.precede(prepare_done);
if (data->prepare_extension.empty()) {
prepare_run.precede(prepare_done);
}
else {
auto prepare_extension = taskflow.compose(
task_prefix + ".extension", data->prepare_extension);
prepare_run.precede(prepare_extension);
prepare_extension.precede(prepare_done);
}
stage_tasks.emplace(root, Stage_Tasks{prepare_if, prepare_done});
}
auto connect_dependencies = [&](const auto& dependency_graph) {
+32 -3
View File
@@ -222,21 +222,31 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity)
auto paint = child.add("paint.visual", [&] { completed.fetch_add(1); });
prepare.precede(paint);
aethera::Task_Graph frame_graph{"test.frame"};
frame_graph.compose("spectrum", child);
auto module_task = frame_graph.compose("spectrum", child);
auto publish = frame_graph.add("publish.state", [&] { completed.fetch_add(1); });
module_task.precede(publish);
aethera::Render_Frame frame{{41, 73}};
frame.mark(aethera::Frame_Trace_Marker::paint_started);
frame.mark(aethera::Frame_Trace_Marker::paint_frame_target_started);
frame.mark(aethera::Frame_Trace_Marker::paint_frame_target_finished);
frame.mark(aethera::Frame_Trace_Marker::paint_finished);
frame.request_taskflow_trace();
ASSERT_TRUE(aethera::detail::begin_taskflow_trace(frame));
aethera::detail::run_taskflow(frame_graph, frame, "test.scene.paint");
aethera::detail::finish_taskflow_trace(frame);
EXPECT_EQ(completed.load(), 2);
EXPECT_EQ(completed.load(), 3);
const auto trace = frame.taskflow_trace();
ASSERT_EQ(trace.graphs.size(), 1u);
EXPECT_EQ(trace.identity, (aethera::Frame_Identity{41, 73}));
EXPECT_TRUE(std::ranges::contains(
trace.markers,
aethera::Frame_Trace_Marker::paint_frame_target_started,
&aethera::Frame_Trace_Point::marker));
EXPECT_EQ(trace.graphs.front().stage, "test.scene.paint");
EXPECT_TRUE(trace.graphs.front().completed);
ASSERT_EQ(trace.graphs.front().nodes.size(), 3u);
ASSERT_EQ(trace.graphs.front().nodes.size(), 4u);
const auto module = std::ranges::find(
trace.graphs.front().nodes, "spectrum",
&aethera::Taskflow_Graph_Trace::Node::name);
@@ -254,4 +264,23 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity)
ASSERT_FALSE(trace.tasks.empty());
for (const auto& task : trace.tasks)
EXPECT_TRUE(metadata_ids.contains(task.native_id));
const auto child_paint = std::ranges::find(
trace.graphs.front().nodes, "paint.visual",
&aethera::Taskflow_Graph_Trace::Node::name);
const auto publish_node = std::ranges::find(
trace.graphs.front().nodes, "publish.state",
&aethera::Taskflow_Graph_Trace::Node::name);
ASSERT_NE(child_paint, trace.graphs.front().nodes.end());
ASSERT_NE(publish_node, trace.graphs.front().nodes.end());
const auto child_paint_execution = std::ranges::find(
trace.tasks, child_paint->native_id,
&aethera::Taskflow_Task_Trace::native_id);
const auto publish_execution = std::ranges::find(
trace.tasks, publish_node->native_id,
&aethera::Taskflow_Task_Trace::native_id);
ASSERT_NE(child_paint_execution, trace.tasks.end());
ASSERT_NE(publish_execution, trace.tasks.end());
EXPECT_NEAR(publish_execution->ready_ms,
child_paint_execution->finished_ms, 0.05);
EXPECT_LT(publish_execution->queue_wait_ms, 1.0);
}