From 2c3105cc0793915cfbb13416cea3f5d25fd2b39a Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Tue, 25 Aug 2026 10:15:59 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8F=90=E5=8D=87=E8=A7=82=E6=B5=8B=E7=B3=BB?= =?UTF-8?q?=E7=BB=9F=20=E9=9A=94=E7=A6=BBtaskflow?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kernel/src/kernel/Task_Graph.cpp | 166 ++++++++++++++ kernel/src/kernel/Task_Graph.hpp | 56 +++++ kernel/src/kernel/Task_Graph_Internal.hpp | 12 + kernel/src/kernel/Taskflow_Frame_Access.hpp | 29 +++ kernel/src/kernel/frame.cpp | 189 ++++++++++++++- kernel/src/kernel/frame.hpp | 44 ++++ kernel/src/kernel/frame_statistics.hpp | 3 + kernel/src/kernel/render_common.cpp | 217 ++++++++++++++++-- kernel/src/kernel/render_common.hpp | 19 +- kernel/src/kernel/render_common.ipp | 10 +- kernel/src/kernel/renderable.cpp | 4 +- kernel/src/kernel/renderable.hpp | 20 +- kernel/src/kernel/renderable.ipp | 24 +- kernel/src/kernel/scene.ipp | 58 ++--- kernel/src/test/render_test.cpp | 54 ++++- render_2D/render_2D/plottable/Afterglow.ipp | 32 +-- .../render_2D/plottable/Frequency_Trace.ipp | 14 +- render_2D/render_2D/plottable/Spectrum.ipp | 18 +- .../render_2D/plottable/Sweep_Spectrum.ipp | 8 +- render_2D/render_2D/plottable/Waterfall.ipp | 24 +- render_2D/render_2D/scene/Render_Scene_2D.cpp | 2 +- render_2D/render_2D/scene/Render_Scene_2D.hpp | 2 +- render_2D/render_2D/scene/Render_Scene_2D.ipp | 65 ++++-- render_2D/tests/Axis_Test.cpp | 8 +- .../render_3D/detail/Async_Render_Backend.cpp | 3 +- render_3D/render_3D/scene/Render_Scene_3D.cpp | 2 +- render_3D/render_3D/scene/Render_Scene_3D.hpp | 2 +- render_3D/render_3D/scene/Render_Scene_3D.ipp | 43 +++- render_3D/render_3D/visual/Basic_Visual.ipp | 10 +- render_3D/tests/Visual_Tests.cpp | 4 +- web_server/src/Gallery_Video_Stream.cpp | 5 +- web_server/src/Plot.cpp | 197 +++++++++++++++- web_server/src/Plot.hpp | 5 + web_server/src/Web_Server.cpp | 101 +++++++- web_server/src/detail/Gallery_Frame_Clock.cpp | 1 + webapp_gallery/src/app.tsx | 215 ++++++++++++++++- webapp_gallery/src/styles.css | 56 +++++ 37 files changed, 1538 insertions(+), 184 deletions(-) create mode 100644 kernel/src/kernel/Task_Graph.cpp create mode 100644 kernel/src/kernel/Task_Graph.hpp create mode 100644 kernel/src/kernel/Task_Graph_Internal.hpp create mode 100644 kernel/src/kernel/Taskflow_Frame_Access.hpp diff --git a/kernel/src/kernel/Task_Graph.cpp b/kernel/src/kernel/Task_Graph.cpp new file mode 100644 index 0000000..ab786b1 --- /dev/null +++ b/kernel/src/kernel/Task_Graph.cpp @@ -0,0 +1,166 @@ +#include "Task_Graph.hpp" +#include "Task_Graph_Internal.hpp" +#include +#include +#include +#include + +namespace aethera { +struct Task_Graph::Private { + struct Node { + tf::Task task{}; /* 原生节点句柄;只在静态构图期间使用。 */ + std::string name{}; /* 调用方提供的业务节点名。 */ + std::string node_id{}; /* 本图内去重后的业务节点 ID。 */ + std::shared_ptr child{}; /* 模块引用的子图;普通节点为空。 */ + }; + explicit Private(std::string graph_name) + : taskflow(std::move(graph_name)) {} + tf::Taskflow taskflow{}; /* 第三方 Taskflow 仅存在于本实现单元。 */ + std::vector nodes{}; /* 与原生 graph 插入顺序一致的业务节点。 */ + std::unordered_map name_counts{}; /* 同名业务节点的下一个序号。 */ + std::uint64_t generation{1}; /* clear/重建后使旧 Task_Node 句柄失效。 */ + + [[nodiscard]] std::string unique_node_id(const std::string& name) { + const std::string base = taskflow.name().empty() + ? name : taskflow.name() + "/" + name; + const auto occurrence = name_counts[name]++; + return occurrence == 0 ? base + : base + "#" + std::to_string(occurrence); + } + void collect_nodes(std::string_view parent, + std::vector& result) const { + for (const auto& source : nodes) { + Taskflow_Graph_Trace::Node node{}; + node.native_id = static_cast(source.task.hash_value()); + auto local_id = std::string_view(source.node_id); + const auto& graph_name = taskflow.name(); + if (!graph_name.empty() && local_id.starts_with(graph_name) && + local_id.size() > graph_name.size() && + local_id[graph_name.size()] == '/') + local_id.remove_prefix(graph_name.size() + 1); + node.node_id = parent.empty() ? source.node_id + : std::string(parent) + "/" + std::string(local_id); + node.parent_node_id = std::string(parent); + node.name = source.name; + node.type = std::string(tf::to_string(source.task.type())); + source.task.for_each_predecessor([&](tf::Task value) { + node.predecessors.push_back( + static_cast(value.hash_value())); + }); + source.task.for_each_successor([&](tf::Task value) { + node.successors.push_back( + static_cast(value.hash_value())); + }); + const auto module_id = node.node_id; + result.push_back(std::move(node)); + if (source.child) source.child->collect_nodes(module_id, result); + } + } +}; + +struct Task_Node::Private { + std::shared_ptr graph{}; /* 节点所属 DAG 及其静态生命周期。 */ + std::size_t index{}; /* DAG 内节点槽位;只在实现单元解释。 */ + std::uint64_t generation{}; /* 防止 clear 后的旧句柄指向复用槽位。 */ +}; + +namespace detail { +void* Task_Graph_Access::native_storage(Task_Graph& graph) noexcept { + return std::addressof(graph.d->taskflow); +} + +std::vector Task_Graph_Access::nodes( + const Task_Graph& graph) { + std::vector result; + graph.d->collect_nodes({}, result); + return result; +} +} + +Task_Node::Task_Node() = default; +Task_Node::~Task_Node() = default; +Task_Node::Task_Node(const Task_Node&) = default; +Task_Node::Task_Node(Task_Node&&) noexcept = default; +Task_Node& Task_Node::operator=(const Task_Node&) = default; +Task_Node& Task_Node::operator=(Task_Node&&) noexcept = default; +Task_Node::Task_Node(std::shared_ptr private_data) + : d(std::move(private_data)) {} + +void Task_Node::precede(const Task_Node& after) const { + if (!d || !after.d || !d->graph || d->graph != after.d->graph || + d->generation != d->graph->generation || + after.d->generation != after.d->graph->generation || + d->index >= d->graph->nodes.size() || + after.d->index >= d->graph->nodes.size()) + throw std::invalid_argument( + "Task graph dependency crosses graph ownership"); + d->graph->nodes[d->index].task.precede( + d->graph->nodes[after.d->index].task); +} + +Task_Graph::Task_Graph(std::string graph_name) + : d(std::make_shared(std::move(graph_name))) {} +Task_Graph::~Task_Graph() = default; +Task_Graph::Task_Graph(Task_Graph&&) noexcept = default; +Task_Graph& Task_Graph::operator=(Task_Graph&& other) noexcept { + if (this == &other) return *this; + if (!d) { + d = std::move(other.d); + return *this; + } + if (!other.d) { + clear(); + return *this; + } + d->taskflow = std::move(other.d->taskflow); + d->nodes = std::move(other.d->nodes); + d->name_counts = std::move(other.d->name_counts); + ++d->generation; + return *this; +} + +Task_Node Task_Graph::add(std::string task_name, + std::function work) { + if (!work) throw std::invalid_argument("Task graph work is empty"); + auto task = d->taskflow.emplace(std::move(work)); + 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}); + return Task_Node{std::make_shared( + Task_Node::Private{d, d->nodes.size() - 1, d->generation})}; +} + +Task_Node Task_Graph::add_condition(std::string task_name, + std::function work) { + if (!work) throw std::invalid_argument("Task graph condition is empty"); + auto task = d->taskflow.emplace(std::move(work)); + 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}); + return Task_Node{std::make_shared( + Task_Node::Private{d, d->nodes.size() - 1, d->generation})}; +} + +Task_Node Task_Graph::compose(std::string task_name, Task_Graph& child) { + auto task = d->taskflow.composed_of(child.d->taskflow); + 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, child.d}); + return Task_Node{std::make_shared( + Task_Node::Private{d, d->nodes.size() - 1, d->generation})}; +} + +void Task_Graph::clear() { + d->taskflow.clear(); + d->nodes.clear(); + d->name_counts.clear(); + ++d->generation; +} +bool Task_Graph::empty() const noexcept { return d->taskflow.empty(); } +std::size_t Task_Graph::size() const noexcept { + return d->taskflow.num_tasks(); +} +const std::string& Task_Graph::name() const noexcept { + return d->taskflow.name(); +} +} diff --git a/kernel/src/kernel/Task_Graph.hpp b/kernel/src/kernel/Task_Graph.hpp new file mode 100644 index 0000000..f4c2425 --- /dev/null +++ b/kernel/src/kernel/Task_Graph.hpp @@ -0,0 +1,56 @@ +#pragma once +#include +#include +#include +#include +namespace aethera { +namespace detail { +struct Task_Graph_Access; +} +class Task_Graph; +/* Task_Graph 中单个业务节点的可复制句柄,仅用于静态构图。 */ +class Task_Node { +public: + Task_Node(); + ~Task_Node(); + Task_Node(const Task_Node&); + Task_Node(Task_Node&&) noexcept; + Task_Node& operator=(const Task_Node&); + Task_Node& operator=(Task_Node&&) noexcept; + /* 建立本节点到 after 的有向依赖;两个节点必须属于同一张图。 */ + void precede(const Task_Node& after) const; +private: + friend class Task_Graph; + struct Private; + explicit Task_Node(std::shared_ptr private_data); + std::shared_ptr d; /* 不透明节点句柄;原生类型只在实现单元可见。 */ +}; +/* + * Kernel 的业务 DAG。Taskflow 类型、Node 句柄和 Observer 绑定全部留在 Implementation 内。 + * 图只允许在未运行时修改;执行期间 add/clear/compose/precede 的行为不受支持。 + */ +class Task_Graph { +public: + explicit Task_Graph(std::string name = {}); + ~Task_Graph(); + Task_Graph(Task_Graph&&) noexcept; + Task_Graph& operator=(Task_Graph&&) noexcept; + Task_Graph(const Task_Graph&) = delete; + Task_Graph& operator=(const Task_Graph&) = delete; + /* 添加普通业务节点;name 会被规范化为帧拓扑中的唯一业务 ID。 */ + Task_Node add(std::string name, std::function work); + /* 添加条件节点;返回的后继下标沿用 Taskflow condition 语义。 */ + Task_Node add_condition(std::string name, std::function work); + /* 添加模块节点并引用 child;child 必须活到本图完成执行。 */ + Task_Node compose(std::string name, Task_Graph& child); + void clear(); + [[nodiscard]] bool empty() const noexcept; + [[nodiscard]] std::size_t size() const noexcept; + [[nodiscard]] const std::string& name() const noexcept; +private: + friend struct Task_Node::Private; + friend struct detail::Task_Graph_Access; + struct Private; + std::shared_ptr d; /* 唯一业务 DAG 定义的不透明所有权。 */ +}; +} diff --git a/kernel/src/kernel/Task_Graph_Internal.hpp b/kernel/src/kernel/Task_Graph_Internal.hpp new file mode 100644 index 0000000..7ed6f20 --- /dev/null +++ b/kernel/src/kernel/Task_Graph_Internal.hpp @@ -0,0 +1,12 @@ +#pragma once +#include "Task_Graph.hpp" +#include "frame.hpp" + +namespace aethera::detail { +/* 仅供 Kernel 实现单元把不透明业务图交给原生 Executor 和帧观测器。 */ +struct Task_Graph_Access { + [[nodiscard]] static void* native_storage(Task_Graph& graph) noexcept; + [[nodiscard]] static std::vector nodes( + const Task_Graph& graph); +}; +} diff --git a/kernel/src/kernel/Taskflow_Frame_Access.hpp b/kernel/src/kernel/Taskflow_Frame_Access.hpp new file mode 100644 index 0000000..56a5c9c --- /dev/null +++ b/kernel/src/kernel/Taskflow_Frame_Access.hpp @@ -0,0 +1,29 @@ +#pragma once +#include "frame.hpp" +#include "Task_Graph.hpp" +#include +#include + +namespace aethera::detail { +struct Taskflow_Graph_Token { + Render_Frame* frame{}; /* DAG 运行记录所属外部帧;帧完成前有效。 */ + std::size_t index{}; /* 帧内 graph 记录索引。 */ +}; + +struct Taskflow_Frame_Access { + using Clock = std::chrono::steady_clock; + static void begin_capture(Render_Frame& frame, std::size_t workers); + static void finish_capture(Render_Frame& frame) noexcept; + [[nodiscard]] static bool acquire_writer(Render_Frame& frame) noexcept; + 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 Taskflow_Graph_Token begin_graph( + Render_Frame& frame, Task_Graph& graph, std::string_view stage); + static void finish_graph(Taskflow_Graph_Token token) noexcept; +}; +} diff --git a/kernel/src/kernel/frame.cpp b/kernel/src/kernel/frame.cpp index 399ee11..45f68ac 100644 --- a/kernel/src/kernel/frame.cpp +++ b/kernel/src/kernel/frame.cpp @@ -1,11 +1,17 @@ #include "frame.hpp" #include "frame_statistics.hpp" +#include "Taskflow_Frame_Access.hpp" +#include "Task_Graph_Internal.hpp" #include #include #include #include #include +#include #include +#include +#include +#include namespace aethera { namespace { constexpr std::size_t marker_count = static_cast(Frame_Trace_Marker::count); @@ -14,11 +20,20 @@ std::uint64_t encode_present_value(std::uint64_t value) noexcept { return value std::uint64_t decode_present_value(std::uint64_t value) noexcept { return value == std::numeric_limits::max() ? value : value - 1; } } struct Render_Frame::Private { + struct Worker_Taskflow_Trace { + std::vector tasks{}; /* 仅对应 worker 写入,捕获结束后统一读取。 */ + }; Frame_Identity identity{}; /* 外部帧管理器提供且终生不变的帧身份。 */ std::chrono::steady_clock::time_point created_at{}; /* 所有 elapsed_ns 使用的单调时钟原点。 */ std::uint64_t created_time_unix_ns{}; /* 用于跨进程展示的创建 Unix 时间,单位为纳秒。 */ std::array markers{}; /* 每种时间点首次出现时的 elapsed_ns 加一编码。 */ std::array measurements{}; /* 每种原始耗时首次记录值的加一编码。 */ + std::atomic_bool taskflow_trace_requested{}; /* 本次逻辑帧是否请求原生 Taskflow 捕获。 */ + std::atomic_bool taskflow_trace_capturing{}; /* 原生 Observer 是否仍可写入本帧。 */ + std::atomic_size_t taskflow_trace_writers{}; /* 正在完成 on_exit 写入的 worker 数。 */ + std::vector taskflow_workers{}; /* 按 Executor worker 隔离的单写者时间线。 */ + std::mutex taskflow_graph_mutex{}; /* 只在按需捕获时保护跨阶段 DAG 追加与完成标记。 */ + std::vector taskflow_graphs{}; /* 本帧主动 run 的业务 DAG 元信息。 */ }; Render_Frame::Render_Frame(Frame_Identity identity) : d(std::make_unique()) { begin(identity); @@ -32,6 +47,14 @@ void Render_Frame::begin(Frame_Identity identity) noexcept { marker.store(0, std::memory_order_relaxed); for (auto& measurement : d->measurements) measurement.store(0, std::memory_order_relaxed); + d->taskflow_trace_requested.store(false, std::memory_order_relaxed); + d->taskflow_trace_capturing.store(false, std::memory_order_relaxed); + d->taskflow_trace_writers.store(0, std::memory_order_relaxed); + d->taskflow_workers.clear(); + { + std::lock_guard lock(d->taskflow_graph_mutex); + d->taskflow_graphs.clear(); + } d->markers[static_cast(Frame_Trace_Marker::created)].store(encode_present_value(0), std::memory_order_relaxed); } Render_Frame::~Render_Frame() = default; @@ -79,6 +102,12 @@ Frame_Statistics_Sample Render_Frame::statistics(Frame_Dimension dimension) cons }; Frame_Statistics_Sample result{}; + result.set(Frame_Statistic::plot_tick_queue_ms, + measurement(Frame_Trace_Measurement::plot_tick_queue_ns)); + result.set(Frame_Statistic::plot_update_ms, + measurement(Frame_Trace_Measurement::plot_update_ns)); + result.set(Frame_Statistic::plot_publish_ms, + measurement(Frame_Trace_Measurement::plot_publish_ns)); result.set(Frame_Statistic::server_completion_ms, marker(Frame_Trace_Marker::frame_ready)); result.set(Frame_Statistic::scene_render_ms, interval( @@ -111,8 +140,19 @@ Frame_Statistics_Sample Render_Frame::statistics(Frame_Dimension dimension) cons Frame_Statistic::gpu_copy_ms, Frame_Statistic::gpu_total_ms, Frame_Statistic::readback_ms}; + constexpr std::array measurement_keys{ + Frame_Trace_Measurement::backend_apply_ns, + Frame_Trace_Measurement::backend_plan_ns, + Frame_Trace_Measurement::backend_execute_ns, + Frame_Trace_Measurement::backend_submit_ns, + Frame_Trace_Measurement::gpu_fence_wait_ns, + Frame_Trace_Measurement::gpu_render_ns, + Frame_Trace_Measurement::gpu_transition_ns, + Frame_Trace_Measurement::gpu_copy_ns, + Frame_Trace_Measurement::gpu_total_ns, + Frame_Trace_Measurement::readback_ns}; for (std::size_t index = 0; index < measurement_statistics.size(); ++index) - result.set(measurement_statistics[index], measurements[index]); + result.set(measurement_statistics[index], measurement(measurement_keys[index])); double remaining = marker(Frame_Trace_Marker::frame_ready); const auto take = [&](double requested) { @@ -201,4 +241,151 @@ Frame_Statistics_Sample Render_Frame::statistics(Frame_Dimension dimension) cons result.set(Frame_Statistic::pipeline_3d_completion_handoff_ms, remaining); return result; } + +void Render_Frame::request_taskflow_trace() noexcept { + d->taskflow_trace_requested.store(true, std::memory_order_release); +} + +bool Render_Frame::taskflow_trace_requested() const noexcept { + return d->taskflow_trace_requested.load(std::memory_order_acquire); +} + +Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { + Taskflow_Frame_Trace result{}; + result.identity = d->identity; + result.created_time_unix_ns = d->created_time_unix_ns; + result.worker_count = d->taskflow_workers.size(); + { + std::lock_guard lock(d->taskflow_graph_mutex); + result.graphs = d->taskflow_graphs; + } + for (const auto& worker : d->taskflow_workers) + result.tasks.insert(result.tasks.end(), worker.tasks.begin(), worker.tasks.end()); + std::ranges::sort(result.tasks, {}, &Taskflow_Task_Trace::started_ms); + + std::erase_if(result.tasks, [&](const Taskflow_Task_Trace& task) { + return std::ranges::none_of(result.graphs, [&](const Taskflow_Graph_Trace& graph) { + return std::ranges::any_of(graph.nodes, [&](const Taskflow_Graph_Trace::Node& node) { + return node.native_id == task.native_id; + }); + }); + }); + + std::unordered_map 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); + } + 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); + } + return result; +} + +void detail::Taskflow_Frame_Access::begin_capture(Render_Frame& frame, std::size_t workers) { + auto& data = *frame.d; + data.taskflow_workers.clear(); + data.taskflow_workers.resize(workers); + for (auto& worker : data.taskflow_workers) worker.tasks.reserve(64); + { + std::lock_guard lock(data.taskflow_graph_mutex); + data.taskflow_graphs.clear(); + } + data.taskflow_trace_writers.store(0, std::memory_order_relaxed); + data.taskflow_trace_capturing.store(true, std::memory_order_release); +} + +void detail::Taskflow_Frame_Access::finish_capture(Render_Frame& frame) noexcept { + auto& data = *frame.d; + data.taskflow_trace_capturing.store(false, std::memory_order_release); + while (data.taskflow_trace_writers.load(std::memory_order_acquire) != 0) + std::this_thread::yield(); +} + +bool detail::Taskflow_Frame_Access::acquire_writer(Render_Frame& frame) noexcept { + auto& data = *frame.d; + if (!data.taskflow_trace_capturing.load(std::memory_order_acquire)) return false; + data.taskflow_trace_writers.fetch_add(1, std::memory_order_acq_rel); + if (data.taskflow_trace_capturing.load(std::memory_order_acquire)) return true; + data.taskflow_trace_writers.fetch_sub(1, std::memory_order_release); + return false; +} + +void detail::Taskflow_Frame_Access::release_writer(Render_Frame& frame) noexcept { + frame.d->taskflow_trace_writers.fetch_sub(1, std::memory_order_release); +} + +bool detail::Taskflow_Frame_Access::contains_task( + const Render_Frame& frame, std::uint64_t native_id) noexcept { + try { + std::lock_guard lock(frame.d->taskflow_graph_mutex); + return std::ranges::any_of( + frame.d->taskflow_graphs, + [&](const Taskflow_Graph_Trace& graph) { + return std::ranges::any_of( + graph.nodes, [&](const Taskflow_Graph_Trace::Node& node) { + return node.native_id == native_id; + }); + }); + } + catch (...) { + return false; + } +} + +void 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) { + auto& data = *frame.d; + if (worker >= data.taskflow_workers.size()) return; + const auto elapsed_ms = [&](Clock::time_point value) { + return std::chrono::duration(value - data.created_at).count(); + }; + Taskflow_Task_Trace trace{}; + trace.native_id = native_id; + trace.worker_id = worker; + trace.worker_queue_size = queue_size; + trace.worker_queue_capacity = queue_capacity; + trace.started_ms = elapsed_ms(started); + trace.finished_ms = elapsed_ms(finished); + trace.duration_ms = std::max(0.0, trace.finished_ms - trace.started_ms); + data.taskflow_workers[worker].tasks.push_back(std::move(trace)); +} + +detail::Taskflow_Graph_Token detail::Taskflow_Frame_Access::begin_graph( + Render_Frame& frame, Task_Graph& taskflow, std::string_view stage) { + Taskflow_Graph_Trace graph{}; + graph.stage = stage; + graph.taskflow_name = taskflow.name(); + graph.nodes = detail::Task_Graph_Access::nodes(taskflow); + graph.submitted_ms = std::chrono::duration( + Clock::now() - frame.d->created_at).count(); + std::lock_guard lock(frame.d->taskflow_graph_mutex); + frame.d->taskflow_graphs.push_back(std::move(graph)); + return {&frame, frame.d->taskflow_graphs.size() - 1}; +} + +void detail::Taskflow_Frame_Access::finish_graph( + detail::Taskflow_Graph_Token token) noexcept { + if (!token.frame) return; + auto& data = *token.frame->d; + std::lock_guard lock(data.taskflow_graph_mutex); + if (token.index >= data.taskflow_graphs.size()) return; + auto& graph = data.taskflow_graphs[token.index]; + graph.finished_ms = std::chrono::duration( + Clock::now() - data.created_at).count(); + graph.completed = true; +} } diff --git a/kernel/src/kernel/frame.hpp b/kernel/src/kernel/frame.hpp index 889565e..52cca73 100644 --- a/kernel/src/kernel/frame.hpp +++ b/kernel/src/kernel/frame.hpp @@ -2,8 +2,10 @@ #include "double_buffer/mechanism.hpp" #include #include +#include #include namespace aethera { +namespace detail { struct Taskflow_Frame_Access; } enum class Frame_Dimension : std::uint8_t; struct Frame_Statistics_Sample; enum class Frame_Trace_Marker : std::uint8_t { @@ -32,6 +34,9 @@ enum class Frame_Trace_Marker : std::uint8_t { count }; enum class Frame_Trace_Measurement : std::uint8_t { + plot_tick_queue_ns, + plot_update_ns, + plot_publish_ns, backend_apply_ns, backend_plan_ns, backend_execute_ns, @@ -57,6 +62,41 @@ struct Frame_Trace_Value { Frame_Trace_Measurement measurement{}; /* 无法表达为公共时间点的原始后端测量类型。 */ std::uint64_t value_ns{}; /* 后端直接记录的持续时间,单位为纳秒。 */ }; +struct Taskflow_Graph_Trace { + std::string stage{}; /* 本帧中执行该原生 Taskflow 的业务阶段。 */ + std::string taskflow_name{}; /* Taskflow 原生图名称。 */ + struct Node { + std::uint64_t native_id{}; /* Taskflow Node 的 native hash。 */ + std::string node_id{}; /* 图路径、业务名和同名序号组成的稳定帧内 ID。 */ + std::string parent_node_id{}; /* 模块子图所属的业务节点;顶层为空。 */ + std::string name{}; /* Taskflow 节点业务名称。 */ + std::string type{}; /* Taskflow 原生 TaskType。 */ + std::vector predecessors{}; /* 原生直接前驱 hash。 */ + std::vector successors{}; /* 原生直接后继 hash。 */ + }; + std::vector nodes{}; /* DAG 构造时登记的节点与依赖元信息。 */ + double submitted_ms{}; /* 相对帧创建时刻的 run 提交时间。 */ + double finished_ms{}; /* 同步返回或异步 topology 完成时间。 */ + bool completed{}; /* 对应 topology 是否已经结束。 */ +}; +struct Taskflow_Task_Trace { + std::uint64_t native_id{}; /* TaskView::hash_value() 返回的原生 Node 身份。 */ + 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 ready_ms{}; /* 前驱完成或根 run 提交后的估算就绪时间。 */ + double queue_wait_ms{}; /* ready 到 on_entry 的估算 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 graphs{}; /* 本帧主动执行的业务 DAG 元信息。 */ + std::vector tasks{}; /* 本帧窗口内原生 Observer 完成的任务执行。 */ +}; class Render_Frame : public double_buffer::Pinned { public: explicit Render_Frame(Frame_Identity identity); @@ -66,11 +106,15 @@ public: void mark(Frame_Trace_Marker marker) noexcept; void record(Frame_Trace_Measurement measurement, std::uint64_t value_ns) noexcept; [[nodiscard]] Frame_Statistics_Sample statistics(Frame_Dimension dimension) const; + void request_taskflow_trace() noexcept; + [[nodiscard]] bool taskflow_trace_requested() const noexcept; + [[nodiscard]] Taskflow_Frame_Trace taskflow_trace() const; protected: /* 物理帧槽再次承载新逻辑帧时,重建其唯一身份和诊断时间原点。 */ void begin(Frame_Identity identity) noexcept; private: struct Private; + friend struct detail::Taskflow_Frame_Access; std::unique_ptr d; /* 帧身份、时钟原点与原子诊断槽位的唯一所有权。 */ }; } diff --git a/kernel/src/kernel/frame_statistics.hpp b/kernel/src/kernel/frame_statistics.hpp index d7d0231..1af50db 100644 --- a/kernel/src/kernel/frame_statistics.hpp +++ b/kernel/src/kernel/frame_statistics.hpp @@ -11,6 +11,9 @@ namespace aethera { enum class Frame_Dimension : std::uint8_t { two_dimensional, three_dimensional }; enum class Frame_Statistic : std::uint8_t { + plot_tick_queue_ms, + plot_update_ms, + plot_publish_ms, server_completion_ms, scene_render_ms, event_dispatch_ms, diff --git a/kernel/src/kernel/render_common.cpp b/kernel/src/kernel/render_common.cpp index 53f8244..098186b 100644 --- a/kernel/src/kernel/render_common.cpp +++ b/kernel/src/kernel/render_common.cpp @@ -1,12 +1,21 @@ #include "render_common.hpp" +#include "Taskflow_Frame_Access.hpp" +#include "Task_Graph_Internal.hpp" #include #include #include #include +#include #include #include +#include namespace aethera { namespace { +tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept { + return *static_cast( + detail::Task_Graph_Access::native_storage(graph)); +} + class Task_Observer : public tf::ObserverInterface { private: using Clock = std::chrono::steady_clock; @@ -18,14 +27,27 @@ private: }; struct Worker_Statistics { std::atomic_size_t task_count{}; + std::atomic_size_t current_queue_size{}; + std::atomic_size_t current_queue_capacity{}; std::atomic_size_t peak_queue_size{}; std::atomic_size_t max_queue_capacity{}; + std::atomic_uint64_t active_task_hash{}; + std::atomic_uint64_t active_task_started_ns{}; + std::atomic active_task_type{tf::TaskType::UNDEFINED}; 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::max()}; std::atomic_uint64_t max_task_time_ns{}; }; - std::vector> starts; + struct Start_Record { + Clock::time_point started{}; /* Observer on_entry 时间。 */ + Render_Frame* frame{}; /* 进入任务时唯一活动的按帧捕获。 */ + std::size_t queue_size{}; /* 进入任务时 worker 队列深度。 */ + std::size_t queue_capacity{}; /* 进入任务时 worker 队列容量。 */ + std::uint64_t native_id{}; /* 嵌套 corun 返回外层任务时恢复其原生身份。 */ + tf::TaskType type{tf::TaskType::UNDEFINED}; /* 嵌套 corun 返回外层任务时恢复其原生类型。 */ + }; + std::vector> starts; std::vector worker_busy_starts; std::unique_ptr worker_statistics; std::size_t worker_statistics_count{}; @@ -49,6 +71,9 @@ private: std::atomic_uint64_t longest_task_time_ns{}; std::atomic_size_t longest_task_hash{}; std::atomic longest_task_type{tf::TaskType::UNDEFINED}; + std::atomic> longest_task_name{}; + std::atomic trace_frame{}; + std::shared_mutex trace_mutex{}; /* 仅按需捕获时保护 Frame* 获取与关闭。 */ static std::uint64_t clock_ns(Clock::time_point value) noexcept { return static_cast(std::chrono::duration_cast(value.time_since_epoch()).count()); } @@ -84,12 +109,37 @@ public: auto active = active_workers.fetch_add(1, std::memory_order_relaxed) + 1; update_max(peak_active_workers, active); } - worker_starts.push_back(now); + worker_starts.push_back(Start_Record{ + now, nullptr, worker.queue_size(), worker.queue_capacity(), + static_cast(task.hash_value()), task.type()}); + auto* frame = trace_frame.load(std::memory_order_acquire); + /* + * 帧租约从 on_entry 持续到对应 on_exit。只在退出时登记写入者会留下 + * “捕获已关闭、物理帧已复用、迟到 on_exit 仍访问旧 Frame*”的窗口。 + */ + if (frame) { + std::shared_lock trace_guard(trace_mutex); + frame = trace_frame.load(std::memory_order_acquire); + if (frame && detail::Taskflow_Frame_Access::acquire_writer(*frame)) { + if (!detail::Taskflow_Frame_Access::contains_task( + *frame, static_cast(task.hash_value()))) { + detail::Taskflow_Frame_Access::release_writer(*frame); + frame = nullptr; + } + } + else frame = nullptr; + } + worker_starts.back().frame = frame; 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()]; + worker_state.current_queue_size.store(worker.queue_size(), std::memory_order_relaxed); + worker_state.current_queue_capacity.store(worker.queue_capacity(), std::memory_order_relaxed); + worker_state.active_task_hash.store(task.hash_value(), std::memory_order_relaxed); + worker_state.active_task_started_ns.store(clock_ns(now), std::memory_order_relaxed); + worker_state.active_task_type.store(task.type(), std::memory_order_relaxed); 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()); @@ -104,7 +154,7 @@ public: auto& worker_starts = starts[worker.id()]; auto start = worker_starts.back(); worker_starts.pop_back(); - auto elapsed = static_cast(std::chrono::duration_cast(now - start).count()); + auto elapsed = static_cast(std::chrono::duration_cast(now - 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()]; @@ -125,6 +175,9 @@ public: 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); + longest_task_name.store( + std::make_shared(task.name()), + std::memory_order_release); } active_tasks.fetch_sub(1, std::memory_order_relaxed); if (worker_starts.empty()) { @@ -132,8 +185,51 @@ public: 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); + worker_state.active_task_hash.store(0, std::memory_order_relaxed); + worker_state.active_task_started_ns.store(0, std::memory_order_relaxed); + worker_state.active_task_type.store(tf::TaskType::UNDEFINED, + std::memory_order_relaxed); + worker_state.current_queue_size.store(0, std::memory_order_relaxed); + worker_state.current_queue_capacity.store(0, std::memory_order_relaxed); + } + 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), + 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(task.hash_value()), + start.queue_size, start.queue_capacity, start.started, now); + } + catch (...) { + /* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */ + } + detail::Taskflow_Frame_Access::release_writer(*start.frame); + } + } + bool begin_trace(Render_Frame& frame, std::size_t workers) { + std::unique_lock guard(trace_mutex); + const auto current = trace_frame.load(std::memory_order_acquire); + if (current) return current == &frame; + detail::Taskflow_Frame_Access::begin_capture(frame, workers); + trace_frame.store(&frame, std::memory_order_release); + return true; + } + void finish_trace(Render_Frame& frame) noexcept { + std::unique_lock guard(trace_mutex); + if (trace_frame.load(std::memory_order_acquire) != &frame) return; + trace_frame.store(nullptr, std::memory_order_release); + detail::Taskflow_Frame_Access::finish_capture(frame); } void write_state(Task_Runtime_State& state, std::size_t workers, std::size_t active_topologies) const { state.worker_count = workers; @@ -156,9 +252,11 @@ public: 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(state.worker_busy_time_ns) * 100.0 / static_cast(state.observed_wall_time_ns) / static_cast(workers) : 0.0; + state.task_types.resize(task_types.size()); for (std::size_t i = 0; i < task_types.size(); ++i) { const auto& source = task_types[i]; auto& target = state.task_types[i]; + target.name = std::string(tf::to_string(tf::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); @@ -166,13 +264,23 @@ public: target.max_time_ns = source.max_time_ns.load(std::memory_order_relaxed); } state.workers.resize(worker_statistics_count); + const auto read_time_ns = clock_ns(Clock::now()); 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.current_queue_size = source.current_queue_size.load(std::memory_order_relaxed); + target.current_queue_capacity = source.current_queue_capacity.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.active_task_hash = source.active_task_hash.load(std::memory_order_relaxed); + const auto active_started = source.active_task_started_ns.load(std::memory_order_relaxed); + target.active_task_time_ns = active_started && read_time_ns >= active_started + ? read_time_ns - active_started : 0; + const auto active_type = source.active_task_type.load(std::memory_order_relaxed); + target.active_task_type = active_started + ? std::string(tf::to_string(active_type)) : std::string{}; 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; @@ -183,8 +291,10 @@ public: } 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); + const auto longest_name = longest_task_name.load(std::memory_order_acquire); + state.longest_task_name = longest_name ? *longest_name : std::string{}; + state.longest_task_type = std::string(tf::to_string( + longest_task_type.load(std::memory_order_relaxed))); } }; class Task_Resource : Pinned { @@ -200,15 +310,22 @@ private: double_buffer::detail::State_Callback_Storage state_callbacks; std::recursive_mutex state_mutex; std::atomic_bool state_callback_enabled{}; + std::atomic_bool executor_ready{}; void create_executor(std::size_t workers, std::shared_ptr worker_interface) { + executor_ready.store(false, std::memory_order_relaxed); executor = std::make_unique(workers, std::move(worker_interface)); - observer.reset(); + observer = executor->make_observer(); state = {}; + executor_ready.store(true, std::memory_order_release); } void ensure_executor() { + /* 正常渲染热路径只读取一次原子位;互斥量仅处理首次惰性初始化。 */ + if (executor_ready.load(std::memory_order_acquire)) return; std::lock_guard guard(state_mutex); if (executor) return; executor = std::make_unique(); + observer = executor->make_observer(); + executor_ready.store(true, std::memory_order_release); } void publish_state() { if (!state_callback_enabled.load(std::memory_order_acquire)) return; @@ -235,7 +352,7 @@ public: completed_taskflows.store(0, std::memory_order_relaxed); failed_taskflows.store(0, std::memory_order_relaxed); } - std::uint64_t run(tf::Taskflow& taskflow) { + std::uint64_t run(Task_Graph& 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); @@ -243,9 +360,9 @@ public: auto start = std::chrono::steady_clock::now(); try { if (executor->this_worker()) - executor->corun(taskflow); + executor->corun(native_taskflow(taskflow)); else - executor->run(taskflow).get(); + executor->run(native_taskflow(taskflow)).get(); } catch (...) { active_taskflows.fetch_sub(1, std::memory_order_relaxed); @@ -259,24 +376,48 @@ public: publish_state(); return elapsed; } - void run(tf::Taskflow& taskflow, std::function completion) { + void run(Task_Graph& taskflow, std::function 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 { + executor->run(native_taskflow(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 task) { + std::uint64_t run(Task_Graph& taskflow, Render_Frame& frame, + std::string_view stage) { + const auto token = detail::Taskflow_Frame_Access::begin_graph( + frame, taskflow, stage); + try { + const auto elapsed = run(taskflow); + detail::Taskflow_Frame_Access::finish_graph(token); + return elapsed; + } + catch (...) { + detail::Taskflow_Frame_Access::finish_graph(token); + throw; + } + } + void run(Task_Graph& taskflow, Render_Frame& frame, std::string_view stage, + std::function completion) { + const auto token = detail::Taskflow_Frame_Access::begin_graph( + frame, taskflow, stage); + run(taskflow, [token, completion = std::move(completion)]() mutable { + detail::Taskflow_Frame_Access::finish_graph(token); + completion(); + }); + } + void schedule(std::string name, std::function task) { if (!task) throw std::invalid_argument("Taskflow scheduled task is empty"); + if (name.empty()) throw std::invalid_argument("Taskflow scheduled task name is empty"); ensure_executor(); - executor->silent_async(std::move(task)); + executor->silent_async(std::move(name), std::move(task)); } std::pmr::memory_resource* memory_resource() const noexcept { return memory; @@ -284,7 +425,6 @@ public: void set_state_callback(std::function callback) { ensure_executor(); std::lock_guard guard(state_mutex); - if (!observer) observer = executor->make_observer(); state_callbacks.template set(std::move(callback)); state_callback_enabled.store(true, std::memory_order_release); } @@ -293,13 +433,34 @@ public: state_callbacks.template clear(); state_callback_enabled.store(false, std::memory_order_release); } + bool begin_trace(Render_Frame& frame) { + ensure_executor(); + return observer->begin_trace(frame, executor->num_workers()); + } + void finish_trace(Render_Frame& frame) noexcept { + if (observer) observer->finish_trace(frame); + } + Task_Runtime_State runtime_state() { + ensure_executor(); + std::lock_guard guard(state_mutex); + Task_Runtime_State result{}; + observer->write_state(result, executor->num_workers(), executor->num_topologies()); + result.active_taskflow_count = active_taskflows.load(std::memory_order_relaxed); + result.peak_active_taskflow_count = peak_active_taskflows.load(std::memory_order_relaxed); + result.completed_taskflow_count = completed_taskflows.load(std::memory_order_relaxed); + result.failed_taskflow_count = failed_taskflows.load(std::memory_order_relaxed); + return result; + } }; } -void initialize_runtime(std::size_t workers, std::shared_ptr worker_interface, Pmr pmr) { - Task_Resource::instance().initialize(workers, std::move(worker_interface), pmr); +void initialize_runtime(std::size_t workers, Pmr pmr) { + Task_Resource::instance().initialize(workers, nullptr, pmr); } -void schedule_task(std::function task) { - Task_Resource::instance().schedule(std::move(task)); +void schedule_task(std::string name, std::function task) { + Task_Resource::instance().schedule(std::move(name), std::move(task)); +} +Task_Runtime_State task_runtime_state() { + return Task_Resource::instance().runtime_state(); } namespace detail { void set_runtime_state_callback_impl(std::function callback) { @@ -308,12 +469,28 @@ void set_runtime_state_callback_impl(std::function completion) { +void run_taskflow(Task_Graph& taskflow, std::function completion) { Task_Resource::instance().run(taskflow, std::move(completion)); } +std::uint64_t run_taskflow(Task_Graph& taskflow, Render_Frame& frame, + std::string_view stage) { + return Task_Resource::instance().run(taskflow, frame, stage); +} +void run_taskflow(Task_Graph& taskflow, Render_Frame& frame, + std::string_view stage, std::function completion) { + Task_Resource::instance().run(taskflow, frame, stage, std::move(completion)); +} +bool begin_taskflow_trace(Render_Frame& frame) { + return !frame.taskflow_trace_requested() || + Task_Resource::instance().begin_trace(frame); +} +void finish_taskflow_trace(Render_Frame& frame) noexcept { + if (frame.taskflow_trace_requested()) + Task_Resource::instance().finish_trace(frame); +} std::pmr::memory_resource* task_memory_resource() noexcept { return Task_Resource::instance().memory_resource(); } diff --git a/kernel/src/kernel/render_common.hpp b/kernel/src/kernel/render_common.hpp index f3102a8..7206799 100644 --- a/kernel/src/kernel/render_common.hpp +++ b/kernel/src/kernel/render_common.hpp @@ -1,6 +1,7 @@ #pragma once #include "double_buffer/model.hpp" #include "frame.hpp" +#include "Task_Graph.hpp" #include #include #include @@ -10,7 +11,6 @@ #include #include #include -#include namespace aethera { using double_buffer::Pinned; using double_buffer::Def; @@ -32,6 +32,7 @@ struct Prepare_Data_Tag {}; struct Task_Runtime_State_Tag {}; /* 单一 Taskflow 任务类型的累计统计。 */ struct Task_Type_State { + std::string name{}; std::size_t count{}; std::uint64_t total_time_ns{}; std::uint64_t min_time_ns{}; @@ -42,8 +43,13 @@ struct Task_Type_State { struct Task_Worker_State { std::size_t id{}; std::size_t task_count{}; + std::size_t current_queue_size{}; /* 当前活跃任务进入时观察到的 Worker 队列深度。 */ + std::size_t current_queue_capacity{}; /* 当前活跃任务进入时观察到的 Worker 队列容量。 */ std::size_t peak_observed_queue_size{}; std::size_t max_observed_queue_capacity{}; + std::uint64_t active_task_hash{}; /* 当前最内层原生任务身份;空闲时为 0。 */ + std::uint64_t active_task_time_ns{}; /* 当前任务从 on_entry 到本次读取已经持续的时间。 */ + std::string active_task_type{}; /* 当前任务的 Taskflow 原生 TaskType;空闲时为空。 */ std::uint64_t task_time_ns{}; std::uint64_t busy_time_ns{}; std::uint64_t idle_time_ns{}; @@ -77,13 +83,13 @@ struct Task_Runtime_State : State_Type { std::size_t max_weak_dependencies{}; std::size_t longest_task_hash{}; std::string longest_task_name; - tf::TaskType longest_task_type{tf::TaskType::UNDEFINED}; + std::string longest_task_type; std::uint64_t longest_task_time_ns{}; std::uint64_t total_task_time_ns{}; std::uint64_t worker_busy_time_ns{}; std::uint64_t observed_wall_time_ns{}; double worker_utilization{}; - std::array task_types{}; + std::vector task_types; std::vector workers; bool operator==(const Task_Runtime_State&) const = default; }; @@ -93,10 +99,11 @@ struct Task_Runtime_State : State_Type { * 未主动调用时运行时会在第一次执行 Taskflow 时按 Taskflow 默认配置惰性创建 Executor。 */ void initialize_runtime(std::size_t workers = std::thread::hardware_concurrency(), - std::shared_ptr worker_interface = nullptr, Pmr pmr = {}); -/* 把独立业务任务提交给全局 Taskflow worker;任务不得执行阻塞式设备等待。 */ -void schedule_task(std::function task); +/* 把具名独立业务任务提交给全局 Taskflow worker;任务不得执行阻塞式设备等待。 */ +void schedule_task(std::string name, std::function task); +/* 读取全局 Executor 的累计状态;计算只读取原子计数,不触发逐任务导出。 */ +[[nodiscard]] Task_Runtime_State task_runtime_state(); /* 为全局 Taskflow 运行时状态注册回调;Tag 目前只接受 Task_Runtime_State_Tag。 */ template Tag, std::invocable Callback> void set_runtime_state_callback(Callback&& callback); diff --git a/kernel/src/kernel/render_common.ipp b/kernel/src/kernel/render_common.ipp index d0c5f7e..91bebe6 100644 --- a/kernel/src/kernel/render_common.ipp +++ b/kernel/src/kernel/render_common.ipp @@ -1,9 +1,15 @@ #pragma once +#include namespace aethera::detail { void set_runtime_state_callback_impl(std::function callback); void clear_runtime_state_callback_impl(); -std::uint64_t run_taskflow(tf::Taskflow& taskflow); -void run_taskflow(tf::Taskflow& taskflow, std::function completion); +std::uint64_t run_taskflow(Task_Graph& taskflow); +void run_taskflow(Task_Graph& taskflow, std::function completion); +std::uint64_t run_taskflow(Task_Graph& taskflow, Render_Frame& frame, std::string_view stage); +void run_taskflow(Task_Graph& taskflow, Render_Frame& frame, std::string_view stage, + std::function completion); +bool begin_taskflow_trace(Render_Frame& frame); +void finish_taskflow_trace(Render_Frame& frame) noexcept; std::pmr::memory_resource* task_memory_resource() noexcept; } namespace aethera { diff --git a/kernel/src/kernel/renderable.cpp b/kernel/src/kernel/renderable.cpp index 67e8e4e..1736b15 100644 --- a/kernel/src/kernel/renderable.cpp +++ b/kernel/src/kernel/renderable.cpp @@ -4,10 +4,10 @@ std::optional Renderable::event_routing_distance(const Event& event) con const auto& data = static_cast(*d); return data.event_routing_distance_run ? data.event_routing_distance_run(this, event) : std::nullopt; } -tf::Taskflow& Renderable::prepare_taskflow() { +Task_Graph& Renderable::prepare_taskflow() { return static_cast(*d).prepare_extension; } -tf::Taskflow& Renderable::paint_taskflow() { +Task_Graph& Renderable::paint_taskflow() { return static_cast(*d).paint_extension; } } diff --git a/kernel/src/kernel/renderable.hpp b/kernel/src/kernel/renderable.hpp index a393140..43c49d2 100644 --- a/kernel/src/kernel/renderable.hpp +++ b/kernel/src/kernel/renderable.hpp @@ -24,12 +24,12 @@ concept Prepare_Data_Renderable = Attached && requires(typename T::Private& p }; /* * Prepare 子图模式定制点。 - * 最终对象的 Private 提供 tf::Taskflow build_prepare_graph(T* object, const T::State& state) 即满足。 + * 最终对象的 Private 提供 Task_Graph build_prepare_graph(T* object, const T::State& state) 即满足。 * 子图首次执行前一定构建,之后由 should_rebuild_prepare_graph(...) 决定是否重建。 */ template concept Prepare_Graph_Renderable = Attached && requires(typename T::Private& private_data, T* object, const typename T::Prop& prop) { - { private_data.build_prepare_graph(object, prop) } -> std::same_as; + { private_data.build_prepare_graph(object, prop) } -> std::same_as; }; /* 最终 Private 声明 No_Prepare 时,该 Renderable 不参与 Prepare 阶段。 */ template @@ -44,12 +44,12 @@ concept Paint_Data_Renderable = Attached && requires(typename T::Private& pri }; /* * Paint 子图模式定制点。 - * 最终对象的 Private 提供 tf::Taskflow build_paint_graph(T* object, const T::State& state) 即满足。 + * 最终对象的 Private 提供 Task_Graph build_paint_graph(T* object, const T::State& state) 即满足。 * 子图首次执行前一定构建,之后由 should_rebuild_paint_graph(...) 决定是否重建。 */ template concept Paint_Graph_Renderable = Attached && requires(typename T::Private& private_data, T* object, const typename T::Prop& prop) { - { private_data.build_paint_graph(object, prop) } -> std::same_as; + { private_data.build_paint_graph(object, prop) } -> std::same_as; }; /* * 最终可交给 Scene 执行的 Renderable 契约。 @@ -89,15 +89,15 @@ struct Renderable : Def { /* 返回该对象参与当前事件竞争时的几何距离;无值表示沿用普通绘制层级路由。 */ [[nodiscard]] std::optional event_routing_distance(const Event& event) const; /* - * 返回 Prepare 阶段完成后执行的直接 Taskflow 扩展端口。 - * tf::Taskflow 只能在该 Renderable 所属 Scene 没有运行时修改;禁止在图执行期间 emplace/erase/clear。 + * 返回 Prepare 阶段完成后执行的业务 DAG 扩展端口。 + * Task_Graph 只能在该 Renderable 所属 Scene 没有运行时修改。 */ - [[nodiscard]] tf::Taskflow& prepare_taskflow(); + [[nodiscard]] Task_Graph& prepare_taskflow(); /* - * 返回 Paint 阶段完成后执行的直接 Taskflow 扩展端口。 - * tf::Taskflow 只能在该 Renderable 所属 Scene 没有运行时修改;禁止在图执行期间 emplace/erase/clear。 + * 返回 Paint 阶段完成后执行的业务 DAG 扩展端口。 + * Task_Graph 只能在该 Renderable 所属 Scene 没有运行时修改。 */ - [[nodiscard]] tf::Taskflow& paint_taskflow(); + [[nodiscard]] Task_Graph& paint_taskflow(); private: /* 数据模式的内部调度入口:执行最终对象 prepare_data(...),再触发各 CRTP 层 after_prepare_data(...)。 */ template diff --git a/kernel/src/kernel/renderable.ipp b/kernel/src/kernel/renderable.ipp index 549c967..973f081 100644 --- a/kernel/src/kernel/renderable.ipp +++ b/kernel/src/kernel/renderable.ipp @@ -1,5 +1,7 @@ #pragma once #include +#include +#include namespace aethera { struct Renderable::Private : Prev_Private { /* @@ -11,7 +13,7 @@ struct Renderable::Private : Prev_Private { using Run_Predicate = bool (*)(Root*, bool); using Rebuild_Predicate = bool (*)(Root*); using Stage_Run = void (*)(Root*); - using Graph_Builder = tf::Taskflow (*)(Root*); + using Graph_Builder = Task_Graph (*)(Root*); using State_Get = State* (*)(Root*); using State_Notify = void (*)(Root*); using Event_Run = void (*)(Root*, const Event&); @@ -33,15 +35,16 @@ struct Renderable::Private : Prev_Private { Stage_Dispatch prepare; /* Prepare 阶段分派。 */ Stage_Dispatch paint; /* Paint 阶段分派。 */ State_Dispatch state; /* Renderable 状态访问与发布分派。 */ + std::string_view business_name; /* 最终 Renderable 类型的诊断业务名。 */ }; const Dispatch* dispatch{}; /* 绑定最终对象类型后指向其静态分派表。 */ Event_Run event_run{}; /* 最终 Private 具备事件能力时的无虚函数入口。 */ Event_Routing_Distance_Run event_routing_distance_run{}; /* 可选事件候选距离;Scene 路由规则按需查询。 */ Color_Cache_Visit color_cache_visit{}; /* 最终对象存在 Color_Cache Buffer 时访问本轮写入结果。 */ - std::unique_ptr prepare_graph; /* Prepare 子图模式的当前构建产物。 */ - std::unique_ptr paint_graph; /* Paint 子图模式的当前构建产物。 */ - tf::Taskflow prepare_extension{}; /* 外部直接续写的 Prepare 完成图;不参与内部子图重建。 */ - tf::Taskflow paint_extension{}; /* 外部直接续写的 Paint 完成图;不参与内部子图重建。 */ + std::unique_ptr prepare_graph; /* Prepare 子图模式的当前构建产物。 */ + std::unique_ptr paint_graph; /* Paint 子图模式的当前构建产物。 */ + Task_Graph prepare_extension{"renderable.prepare.extension"}; /* 外部直接续写的 Prepare 完成图。 */ + Task_Graph paint_extension{"renderable.paint.extension"}; /* 外部直接续写的 Paint 完成图。 */ bool prepare_graph_built{}; /* Prepare 子图是否至少成功构建过一次。 */ bool paint_graph_built{}; /* Paint 子图是否至少成功构建过一次。 */ /* CRTP 可覆盖:决定已选中子图模式的 Prepare 子图是否重建;object 为最终对象,state 为当前发布状态;默认返回 false。 */ @@ -112,6 +115,14 @@ inline void Renderable::bind_dependency_graph_object(Attached auto* object) { using Object = std::remove_pointer_t; static_assert(Renderable_Object); auto& data = static_cast(*object->d); + static const std::string business_name = [] { + std::string name = typeid(typename Object::Attached_Object).name(); + for (const std::string_view prefix : {"struct ", "class "}) + if (name.starts_with(prefix)) name.erase(0, prefix.size()); + if (const auto separator = name.rfind("::"); separator != std::string::npos) + name.erase(0, separator + 2); + return name; + }(); static const Private::Dispatch dispatch{ { [](Root* root, bool dirty) { @@ -202,7 +213,8 @@ inline void Renderable::bind_dependency_graph_object(Attached auto* object) { private_data.state.advance(); value->template notify_state(); } - } + }, + business_name }; data.dispatch = &dispatch; object->template mark_dirty(); diff --git a/kernel/src/kernel/scene.ipp b/kernel/src/kernel/scene.ipp index 42d3cf0..e36bc90 100644 --- a/kernel/src/kernel/scene.ipp +++ b/kernel/src/kernel/scene.ipp @@ -1,4 +1,5 @@ #pragma once +#include "Task_Graph_Internal.hpp" #include #include #include @@ -32,7 +33,7 @@ struct Scene::Private : Prev_Private { /* 派生 Scene 的 Private 还可覆盖 Def::Private 的四个 State 生命周期 hook,并通过 State_Access::get() 访问状态层。 */ }; struct Scene::Private::Runtime { - std::unique_ptr taskflow; /* 当前已构建的总 Taskflow;为空表示尚未构建。 */ + std::unique_ptr taskflow; /* 当前已构建的总业务 DAG;为空表示尚未构建。 */ }; template void Scene::Private::bind_private_crtp(Object* object) { @@ -67,7 +68,10 @@ void Scene::Private::process(Object* object, Render_Frame* frame, Callback&& cal auto& state = static_cast(*private_data.state.pending); state.taskflow_execution_time_ns = 0; if (frame) frame->mark(Frame_Trace_Marker::prepare_started); - if (runtime->taskflow && !runtime->taskflow->empty()) state.taskflow_execution_time_ns = detail::run_taskflow(*runtime->taskflow); + if (runtime->taskflow && !runtime->taskflow->empty()) + state.taskflow_execution_time_ns = frame && frame->taskflow_trace_requested() + ? detail::run_taskflow(*runtime->taskflow, *frame, "scene.prepare") + : detail::run_taskflow(*runtime->taskflow); if (frame) frame->mark(Frame_Trace_Marker::prepare_finished); state.event_statistics = event_statistics.state(); private_data.state.advance(); @@ -102,7 +106,7 @@ void Scene::Private::after_advance(Object* object, object->template access_pending_dependency_graph( [&](auto& prepare_state) { taskflow_dirty = taskflow_dirty || prepare_state.dirty(); }); if (!taskflow_dirty) return; - if (!runtime->taskflow) runtime->taskflow = std::make_unique(); + if (!runtime->taskflow) runtime->taskflow = std::make_unique("scene.prepare"); auto& taskflow = *runtime->taskflow; std::pmr::unordered_set renderables{resource}; prepare_dependencies.for_each_bound( @@ -113,8 +117,8 @@ void Scene::Private::after_advance(Object* object, scene_state.renderable_count = renderables.size(); taskflow.clear(); struct Stage_Tasks { - tf::Task prepare_entry; /* Prepare 条件任务,作为该阶段依赖入口。 */ - tf::Task prepare_exit; /* Prepare 完成任务,作为该阶段依赖出口。 */ + Task_Node prepare_entry; /* Prepare 条件任务,作为该阶段依赖入口。 */ + Task_Node prepare_exit; /* Prepare 完成任务,作为该阶段依赖出口。 */ }; std::pmr::unordered_map stage_tasks{resource}; for (auto* renderable : renderables) { @@ -122,7 +126,8 @@ void Scene::Private::after_advance(Object* object, auto* data = prepare_dependencies.private_data(root); if (!data) continue; auto* dispatch = data->dispatch; - auto prepare_if = taskflow.emplace([data, dispatch, root] { + const auto task_prefix = std::string(dispatch->business_name) + ".prepare"; + auto prepare_if = taskflow.add_condition(task_prefix + ".condition", [data, dispatch, root] { auto& state = *dispatch->state.pending(root); state.prepare_graph_rebuilt = false; state.prepare_execution_time_ns = 0; @@ -134,7 +139,7 @@ void Scene::Private::after_advance(Object* object, state.prepare_graph_rebuilt = true; root->template mark_dirty(); } - state.prepare_task_count = data->prepare_graph->num_tasks(); + state.prepare_task_count = data->prepare_graph->size(); } else { state.prepare_task_count = dispatch->prepare.run ? 1 : 0; @@ -148,20 +153,20 @@ void Scene::Private::after_advance(Object* object, std::chrono::steady_clock::now().time_since_epoch()).count()); } return state.prepare_executed ? 0 : 1; - }).name("renderable.prepare.condition"); - tf::Task prepare_run; + }); + Task_Node prepare_run; if (dispatch->prepare.builder) { - if (!data->prepare_graph) data->prepare_graph = std::make_unique(); - prepare_run = taskflow.composed_of(*data->prepare_graph).name("renderable.prepare.graph"); + if (!data->prepare_graph) data->prepare_graph = std::make_unique(task_prefix + ".graph"); + prepare_run = taskflow.compose(task_prefix + ".graph", *data->prepare_graph); } else { - prepare_run = taskflow.emplace([dispatch, root] { + prepare_run = taskflow.add(task_prefix + ".data", [dispatch, root] { if (dispatch->prepare.run) dispatch->prepare.run(root); - }).name("renderable.prepare.data"); + }); } - auto prepare_extension = taskflow.composed_of(data->prepare_extension) - .name("renderable.prepare.extension"); - auto prepare_done = taskflow.emplace([dispatch, 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) { root->template take_dirty(); @@ -171,8 +176,9 @@ void Scene::Private::after_advance(Object* object, state.prepare_execution_time_ns = finished - state.prepare_execution_time_ns; } dispatch->state.publish(root); - }).name("renderable.prepare.complete"); - prepare_if.precede(prepare_run, prepare_done); + }); + prepare_if.precede(prepare_run); + prepare_if.precede(prepare_done); prepare_run.precede(prepare_extension); prepare_extension.precede(prepare_done); stage_tasks.emplace(root, Stage_Tasks{prepare_if, prepare_done}); @@ -203,17 +209,17 @@ void Scene::Private::after_advance(Object* object, ); }; connect_dependencies(prepare_dependencies); - scene_state.taskflow_task_count = taskflow.num_tasks(); + scene_state.taskflow_task_count = taskflow.size(); scene_state.taskflow_dependency_count = 0; scene_state.taskflow_max_predecessors = 0; scene_state.taskflow_max_successors = 0; - taskflow.for_each_task( - [&](tf::Task task) { - scene_state.taskflow_dependency_count += task.num_successors(); - scene_state.taskflow_max_predecessors = std::max(scene_state.taskflow_max_predecessors, task.num_predecessors()); - scene_state.taskflow_max_successors = std::max(scene_state.taskflow_max_successors, task.num_successors()); - } - ); + for (const auto& node : detail::Task_Graph_Access::nodes(taskflow)) { + scene_state.taskflow_dependency_count += node.successors.size(); + scene_state.taskflow_max_predecessors = std::max( + scene_state.taskflow_max_predecessors, node.predecessors.size()); + scene_state.taskflow_max_successors = std::max( + scene_state.taskflow_max_successors, node.successors.size()); + } object->template access_pending_dependency_graph( [](auto& prepare_state) { if (prepare_state.dirty()) prepare_state.take_dirty(); }); scene_state.taskflow_rebuilt = true; diff --git a/kernel/src/test/render_test.cpp b/kernel/src/test/render_test.cpp index 64eb83c..ebdfae2 100644 --- a/kernel/src/test/render_test.cpp +++ b/kernel/src/test/render_test.cpp @@ -1,5 +1,7 @@ #include "scene.hpp" #include +#include +#include namespace { struct Direct_Renderable : double_buffer::Def { struct Prop : Prev_Prop {}; @@ -28,10 +30,10 @@ struct Graph_Renderable : double_buffer::Def(); add_renderable(*scene, renderable.get()); int prepare_extension_calls{}; - renderable->prepare_taskflow().emplace( - [&] { ++prepare_extension_calls; }).name("test.prepare.extension"); + renderable->prepare_taskflow().add( + "test.prepare.extension", [&] { ++prepare_extension_calls; }); scene->process([](const auto&) {}); auto& data = renderable->data_for_test(); auto& base = static_cast(data); @@ -211,3 +213,45 @@ TEST(scene_condition, upstream_change_makes_downstream_run_in_same_taskflow) { EXPECT_EQ(source_data.prepare_calls, 2); EXPECT_EQ(target_data.prepare_calls, 2); } + +TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) { + aethera::initialize_runtime(2); + std::atomic_int completed{}; + aethera::Task_Graph child{"test.visual"}; + auto prepare = child.add("prepare.samples", [&] { completed.fetch_add(1); }); + 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); + + aethera::Render_Frame frame{{41, 73}}; + 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); + const auto trace = frame.taskflow_trace(); + ASSERT_EQ(trace.graphs.size(), 1u); + EXPECT_EQ(trace.identity, (aethera::Frame_Identity{41, 73})); + EXPECT_EQ(trace.graphs.front().stage, "test.scene.paint"); + EXPECT_TRUE(trace.graphs.front().completed); + ASSERT_EQ(trace.graphs.front().nodes.size(), 3u); + const auto module = std::ranges::find( + trace.graphs.front().nodes, "spectrum", + &aethera::Taskflow_Graph_Trace::Node::name); + ASSERT_NE(module, trace.graphs.front().nodes.end()); + const auto child_prepare = std::ranges::find( + trace.graphs.front().nodes, "prepare.samples", + &aethera::Taskflow_Graph_Trace::Node::name); + ASSERT_NE(child_prepare, trace.graphs.front().nodes.end()); + EXPECT_EQ(child_prepare->parent_node_id, module->node_id); + EXPECT_EQ(child_prepare->node_id, "test.frame/spectrum/prepare.samples"); + + std::unordered_set metadata_ids; + for (const auto& node : trace.graphs.front().nodes) + metadata_ids.insert(node.native_id); + ASSERT_FALSE(trace.tasks.empty()); + for (const auto& task : trace.tasks) + EXPECT_TRUE(metadata_ids.contains(task.native_id)); +} diff --git a/render_2D/render_2D/plottable/Afterglow.ipp b/render_2D/render_2D/plottable/Afterglow.ipp index 5fdeeec..29fd075 100644 --- a/render_2D/render_2D/plottable/Afterglow.ipp +++ b/render_2D/render_2D/plottable/Afterglow.ipp @@ -21,10 +21,10 @@ struct Afterglow::Private : Prev_Private { void bind_sources(Frequency_Object* frequency_axis_value, Power_Object* power_axis_value); /* CRTP 覆盖:构建累加、归一化和着色三阶段 Prepare 子图。 */ template - [[nodiscard]] tf::Taskflow build_prepare_graph(Object* object, const Prop& state); + [[nodiscard]] Task_Graph build_prepare_graph(Object* object, const Prop& state); /* CRTP 覆盖:构建消费色块矩阵的 Paint 子图。 */ template - [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& state); + [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& state); /* CRTP 覆盖:分块数量改变时请求重建 Prepare 子图。 */ template [[nodiscard]] bool should_rebuild_prepare_graph(Object* object, const Prop& state); @@ -85,27 +85,27 @@ bool Afterglow::Private::should_rebuild_prepare_graph(Object* object, const Prop return graph_partition_count != detail::curve_partition_count(state.partition_mode, state.partition_count, std::max(1, columns)); } template -tf::Taskflow Afterglow::Private::build_prepare_graph(Object* object, const Prop& state) { +Task_Graph Afterglow::Private::build_prepare_graph(Object* object, const Prop& state) { const std::size_t available = object->template access_query_stream( [](std::span>> spectra) { return spectra.empty() ? 0 : spectra.back()->size(); }); const std::size_t columns = state.frequency_point_size ? std::min(state.frequency_point_size, available) : available; graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, std::max(1, columns)); - tf::Taskflow graph; - auto begin = graph.emplace([this, object] { + Task_Graph graph{"afterglow.prepare"}; + auto begin = graph.add("frame", [this, object] { prepare_frame(object); - }).name("afterglow.prepare.frame"); - auto normalize = graph.emplace([this] { + }); + auto normalize = graph.add("normalize", [this] { normalize_frame(); - }).name("afterglow.prepare.normalize"); + }); for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { - auto accumulate = graph.emplace([this, object, index] { + auto accumulate = graph.add("accumulate", [this, object, index] { accumulate_partition(object, index); - }).name("afterglow.prepare.accumulate"); - auto color = graph.emplace([this, object, index] { + }); + auto color = graph.add("color", [this, object, index] { color_partition(object, index); - }).name("afterglow.prepare.color"); + }); begin.precede(accumulate); accumulate.precede(normalize); normalize.precede(color); @@ -113,11 +113,11 @@ tf::Taskflow Afterglow::Private::build_prepare_graph(Object* object, const Prop& return graph; } template -tf::Taskflow Afterglow::Private::build_paint_graph(Object* object, const Prop&) { - tf::Taskflow graph; - graph.emplace([this, object] { +Task_Graph Afterglow::Private::build_paint_graph(Object* object, const Prop&) { + Task_Graph graph{"afterglow.paint"}; + graph.add("frame", [this, object] { paint_frame(object); - }).name("afterglow.paint.frame"); + }); return graph; } template diff --git a/render_2D/render_2D/plottable/Frequency_Trace.ipp b/render_2D/render_2D/plottable/Frequency_Trace.ipp index 2da407c..3fb6417 100644 --- a/render_2D/render_2D/plottable/Frequency_Trace.ipp +++ b/render_2D/render_2D/plottable/Frequency_Trace.ipp @@ -15,9 +15,9 @@ struct Frequency_Trace::Private : Prev_Private { Plot_Partition_Count graph_partition_count{}; /* 当前 Prepare 子图固化的分块数。 */ void bind_sources(Time_Object* time_axis_value, Value_Object* value_axis_value); /* CRTP 覆盖:按当前样本规模构建分块 Prepare 子图。 */ - template [[nodiscard]] tf::Taskflow build_prepare_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_prepare_graph(Object* object, const Prop& state); /* CRTP 覆盖:构建消费已准备曲线的 Paint 子图。 */ - template [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& state); /* CRTP 覆盖:分块数量改变时请求重建 Prepare 子图。 */ template [[nodiscard]] bool should_rebuild_prepare_graph(Object* object, const Prop& state); template void prepare_frame(Object* object, Plot_Partition_Count partition_count); @@ -36,14 +36,14 @@ std::expected, Dependency_Graph_Error> Frequency_Trace:: template bool Frequency_Trace::Private::should_rebuild_prepare_graph(Object* object, const Prop& state) { return graph_partition_count != detail::curve_partition_count(state.partition_mode, state.partition_count, object->template access_query_stream([](std::span samples) { return samples.size(); })); } template -tf::Taskflow Frequency_Trace::Private::build_prepare_graph(Object* object, const Prop& state) { - graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, object->template access_query_stream([](std::span samples) { return samples.size(); })); tf::Taskflow graph; - auto begin = graph.emplace([this, object] { prepare_frame(object, graph_partition_count); }).name("frequency_trace.prepare.frame"); - for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { auto task = graph.emplace([this, object, index] { prepare_partition(object, index); }).name("frequency_trace.prepare.partition"); begin.precede(task); } +Task_Graph Frequency_Trace::Private::build_prepare_graph(Object* object, const Prop& state) { + graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, object->template access_query_stream([](std::span samples) { return samples.size(); })); Task_Graph graph{"frequency_trace.prepare"}; + auto begin = graph.add("frame", [this, object] { prepare_frame(object, graph_partition_count); }); + for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { auto task = graph.add("partition", [this, object, index] { prepare_partition(object, index); }); begin.precede(task); } return graph; } template -tf::Taskflow Frequency_Trace::Private::build_paint_graph(Object* object, const Prop&) { tf::Taskflow graph; graph.emplace([this, object] { paint_frame(object); }).name("frequency_trace.paint.frame"); return graph; } +Task_Graph Frequency_Trace::Private::build_paint_graph(Object* object, const Prop&) { Task_Graph graph{"frequency_trace.paint"}; graph.add("frame", [this, object] { paint_frame(object); }); return graph; } template void Frequency_Trace::Private::prepare_frame(Object* object, Plot_Partition_Count partition_count) { const auto& state = object->template read_prop(); const auto& time_layout = time_axis->template read_prop(); const auto& value_layout = value_axis->template read_prop(); diff --git a/render_2D/render_2D/plottable/Spectrum.ipp b/render_2D/render_2D/plottable/Spectrum.ipp index fe63f78..aeb255e 100644 --- a/render_2D/render_2D/plottable/Spectrum.ipp +++ b/render_2D/render_2D/plottable/Spectrum.ipp @@ -42,9 +42,9 @@ struct Spectrum::Private : Prev_Private { /* CRTP State 钩子:Spectrum 业务状态写入后标记自身 Prepare;其他继承层状态由各自 Private 负责。 */ template void after_prop_set(Object* object, Member Owner::* member, Prop_Access pending_states); /* CRTP 子图能力:按当前样本数和 State 分块策略构建并行 Prepare 图。 */ - template [[nodiscard]] tf::Taskflow build_prepare_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_prepare_graph(Object* object, const Prop& state); /* CRTP 子图能力:构建背景、分块曲线和覆盖标记的 Paint 图。 */ - template [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& state); /* CRTP 覆盖:样本规模或分块配置改变时重建 Prepare 子图。 */ template [[nodiscard]] bool should_rebuild_prepare_graph(Object* object, const Prop& state); template [[nodiscard]] std::size_t desired_partition_count(const Object* object, const Prop& state) const; @@ -95,21 +95,21 @@ bool Spectrum::Private::should_rebuild_prepare_graph(Object* object, const Prop& return prepare_graph_partition_count != desired_partition_count(object, state); } template -tf::Taskflow Spectrum::Private::build_prepare_graph(Object* object, const Prop& state) { +Task_Graph Spectrum::Private::build_prepare_graph(Object* object, const Prop& state) { const std::size_t partition_count = desired_partition_count(object, state); prepare_graph_partition_count = partition_count; - tf::Taskflow graph; - auto begin = graph.emplace([this, object, partition_count] { prepare_frame(object, partition_count); }).name("spectrum.prepare.frame"); + Task_Graph graph{"spectrum.prepare"}; + auto begin = graph.add("frame", [this, object, partition_count] { prepare_frame(object, partition_count); }); for (std::size_t index = 0; index < partition_count; ++index) { - auto partition = graph.emplace([this, object, index] { prepare_partition(object, index); }).name("spectrum.prepare.partition"); + auto partition = graph.add("partition", [this, object, index] { prepare_partition(object, index); }); begin.precede(partition); } return graph; } template -tf::Taskflow Spectrum::Private::build_paint_graph(Object* object, const Prop&) { - tf::Taskflow graph; - graph.emplace([this, object] { paint_frame(object); }).name("spectrum.paint.frame"); +Task_Graph Spectrum::Private::build_paint_graph(Object* object, const Prop&) { + Task_Graph graph{"spectrum.paint"}; + graph.add("frame", [this, object] { paint_frame(object); }); return graph; } template diff --git a/render_2D/render_2D/plottable/Sweep_Spectrum.ipp b/render_2D/render_2D/plottable/Sweep_Spectrum.ipp index 91299f4..232c32c 100644 --- a/render_2D/render_2D/plottable/Sweep_Spectrum.ipp +++ b/render_2D/render_2D/plottable/Sweep_Spectrum.ipp @@ -27,9 +27,9 @@ struct Sweep_Spectrum::Private : Prev_Private { Plot_Partition_Count graph_partition_count{}; /* 当前 Prepare 子图分块数。 */ void bind_sources(Frequency_Object* frequency_axis_value, Power_Object* power_axis_value); /* CRTP 覆盖:按扫描点规模构建分块 Prepare 子图。 */ - template [[nodiscard]] tf::Taskflow build_prepare_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_prepare_graph(Object* object, const Prop& state); /* CRTP 覆盖:构建消费曲线分块的 Paint 子图。 */ - template [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& state); + template [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& state); /* CRTP 覆盖:分块数量改变时请求重建 Prepare 子图。 */ template [[nodiscard]] bool should_rebuild_prepare_graph(Object* object, const Prop& state); template void prepare_frame(Object* object); @@ -48,9 +48,9 @@ std::expected, Dependency_Graph_Error> Sweep_Spectrum::B template bool Sweep_Spectrum::Private::should_rebuild_prepare_graph(Object*, const Prop& state) { return graph_partition_count != detail::curve_partition_count(state.partition_mode, state.partition_count, stored_point_count()); } template -tf::Taskflow Sweep_Spectrum::Private::build_prepare_graph(Object* object, const Prop& state) { graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, stored_point_count()); tf::Taskflow graph; auto begin = graph.emplace([this, object] { prepare_frame(object); }).name("sweep_spectrum.prepare.frame"); for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { auto task = graph.emplace([this, object, index] { prepare_partition(object, index); }).name("sweep_spectrum.prepare.partition"); begin.precede(task); } return graph; } +Task_Graph Sweep_Spectrum::Private::build_prepare_graph(Object* object, const Prop& state) { graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, stored_point_count()); Task_Graph graph{"sweep_spectrum.prepare"}; auto begin = graph.add("frame", [this, object] { prepare_frame(object); }); for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { auto task = graph.add("partition", [this, object, index] { prepare_partition(object, index); }); begin.precede(task); } return graph; } template -tf::Taskflow Sweep_Spectrum::Private::build_paint_graph(Object* object, const Prop&) { tf::Taskflow graph; graph.emplace([this, object] { paint_frame(object); }).name("sweep_spectrum.paint.frame"); return graph; } +Task_Graph Sweep_Spectrum::Private::build_paint_graph(Object* object, const Prop&) { Task_Graph graph{"sweep_spectrum.paint"}; graph.add("frame", [this, object] { paint_frame(object); }); return graph; } template void Sweep_Spectrum::Private::prepare_frame(Object* object) { const auto& state = object->template read_prop(); const auto& frequency_layout = frequency_axis->template read_prop(); const auto& power_layout = power_axis->template read_prop(); const auto& power_state = power_axis->template read_prop(); prepared = {}; prepared.canvas = scene->template read_prop().viewport; diff --git a/render_2D/render_2D/plottable/Waterfall.ipp b/render_2D/render_2D/plottable/Waterfall.ipp index 908cf05..1e5f860 100644 --- a/render_2D/render_2D/plottable/Waterfall.ipp +++ b/render_2D/render_2D/plottable/Waterfall.ipp @@ -32,10 +32,10 @@ struct Waterfall::Private : Prev_Private { void bind_sources(Frequency_Object* frequency_axis_value, Time_Object* time_axis_value); /* CRTP 覆盖:按当前色块工作量构建分块 Prepare 子图。 */ template - [[nodiscard]] tf::Taskflow build_prepare_graph(Object* object, const Prop& state); + [[nodiscard]] Task_Graph build_prepare_graph(Object* object, const Prop& state); /* CRTP 覆盖:构建消费色块矩阵和提示信息的 Paint 子图。 */ template - [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& state); + [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& state); /* CRTP 覆盖:分块数量改变时请求重建 Prepare 子图。 */ template [[nodiscard]] bool should_rebuild_prepare_graph(Object* object, const Prop& state); @@ -96,29 +96,29 @@ bool Waterfall::Private::should_rebuild_prepare_graph(Object* object, const Prop return graph_partition_count != detail::curve_partition_count(state.partition_mode, state.partition_count, cells); } template -tf::Taskflow Waterfall::Private::build_prepare_graph(Object* object, const Prop& state) { +Task_Graph Waterfall::Private::build_prepare_graph(Object* object, const Prop& state) { const std::size_t cells = object->template access_query_stream([&](std::span> rows) { return rows.size() * (state.frequency_bin_count ? state.frequency_bin_count : rows.empty() ? 1 : rows.back()->values.size()); }); graph_partition_count = detail::curve_partition_count(state.partition_mode, state.partition_count, cells); - tf::Taskflow graph; - auto begin = graph.emplace([this, object] { + Task_Graph graph{"waterfall.prepare"}; + auto begin = graph.add("frame", [this, object] { prepare_frame(object); - }).name("waterfall.prepare.frame"); + }); for (Plot_Partition_Count index = 0; index < graph_partition_count; ++index) { - auto task = graph.emplace([this, object, index] { + auto task = graph.add("partition", [this, object, index] { prepare_partition(object, index); - }).name("waterfall.prepare.partition"); + }); begin.precede(task); } return graph; } template -tf::Taskflow Waterfall::Private::build_paint_graph(Object* object, const Prop&) { - tf::Taskflow graph; - graph.emplace([this, object] { +Task_Graph Waterfall::Private::build_paint_graph(Object* object, const Prop&) { + Task_Graph graph{"waterfall.paint"}; + graph.add("frame", [this, object] { paint_frame(object); - }).name("waterfall.paint.frame"); + }); return graph; } template diff --git a/render_2D/render_2D/scene/Render_Scene_2D.cpp b/render_2D/render_2D/scene/Render_Scene_2D.cpp index 9d79a88..7fddf19 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.cpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.cpp @@ -9,7 +9,7 @@ Render_Scene_2D::render(Frame_2D* frame) { return static_cast(*d).dispatch->render(this, frame); } void Render_Scene_2D::set_frame_callback(Frame_Callback callback) { static_cast(*d).dispatch->set_frame_callback(this, std::move(callback)); } -tf::Taskflow& Render_Scene_2D::completion_taskflow() { +Task_Graph& Render_Scene_2D::completion_taskflow() { return static_cast(*d).completion_graph; } void Render_Scene_2D::activate_view() { diff --git a/render_2D/render_2D/scene/Render_Scene_2D.hpp b/render_2D/render_2D/scene/Render_Scene_2D.hpp index 945e7ad..8c92019 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.hpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.hpp @@ -48,7 +48,7 @@ struct Render_Scene_2D : Def paint_taskflow{}; /* 仅由二维 Paint 图构建的执行图。 */ + Task_Graph completion_graph{"render_2d.completion"}; /* 最终像素完成后、发布回调前执行的外部续写图。 */ + std::unique_ptr paint_taskflow{}; /* 仅由二维 Paint 图构建的执行图。 */ ~Private(); Frame_2D* active_frame{}; /* 当前同步 process 借用的外部帧;render 返回前清空。 */ Blend2D_Cache* frame_target{}; /* 当前 render(frame) 所属外部颜色层;调用返回后清空。 */ @@ -104,7 +104,7 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access rebuild = true; } if (!rebuild) return; - if (!paint_taskflow) paint_taskflow = std::make_unique(); + if (!paint_taskflow) paint_taskflow = std::make_unique("render_2d.paint"); auto& taskflow = *paint_taskflow; taskflow.clear(); const auto dependencies = object->template current_dependency_graph(); @@ -139,14 +139,15 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access std::unordered_set closed_groups; for (auto current = paint_order.rbegin(); current != paint_order.rend(); ++current) if (current->cache_owner && closed_groups.insert(current->cache_owner).second) current->cache_group_last = true; - tf::Task previous; + Task_Node previous; bool has_previous{}; for (std::size_t index = 0; index < paint_order.size(); ++index) { auto& paint_node = paint_order[index]; Root* root = paint_node.object; auto& data = *paint_node.private_data; auto* dispatch = data.dispatch; - auto paint_if = taskflow.emplace([this, &data, dispatch, root, index] { + const auto task_prefix = std::string(dispatch->business_name) + ".paint"; + auto paint_if = taskflow.add_condition(task_prefix + ".condition", [this, &data, dispatch, root, index] { auto& state = *dispatch->state.pending(root); state.paint_graph_rebuilt = false; state.paint_execution_time_ns = 0; @@ -158,7 +159,7 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access state.paint_graph_rebuilt = true; root->template mark_dirty(); } - state.paint_task_count = data.paint_graph->num_tasks(); + state.paint_task_count = data.paint_graph->size(); } else state.paint_task_count = dispatch->paint.run ? 1 : 0; const bool dirty = root->template dirty(); state.paint_dirty = dirty; @@ -169,15 +170,14 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access std::chrono::steady_clock::now().time_since_epoch()).count()); } return state.paint_executed ? 0 : 1; - }).name("render_2d.paint.condition"); - tf::Task paint_run; + }); + Task_Node paint_run; if (dispatch->paint.builder) { - if (!data.paint_graph) data.paint_graph = std::make_unique(); - paint_run = taskflow.composed_of(*data.paint_graph).name("render_2d.paint.graph"); - } else paint_run = taskflow.emplace([dispatch, root] { if (dispatch->paint.run) dispatch->paint.run(root); }).name("render_2d.paint.data"); - auto paint_extension = taskflow.composed_of(data.paint_extension) - .name("renderable.paint.extension"); - auto paint_done = taskflow.emplace([dispatch, root] { + if (!data.paint_graph) data.paint_graph = std::make_unique(task_prefix + ".graph"); + paint_run = taskflow.compose(task_prefix + ".graph", *data.paint_graph); + } else paint_run = taskflow.add(task_prefix + ".data", [dispatch, root] { if (dispatch->paint.run) dispatch->paint.run(root); }); + auto paint_extension = taskflow.compose(task_prefix + ".extension", data.paint_extension); + auto paint_done = taskflow.add(task_prefix + ".complete", [dispatch, root] { auto& state = *dispatch->state.pending(root); if (state.paint_executed) { root->template take_dirty(); @@ -187,15 +187,16 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access state.paint_execution_time_ns = finished - state.paint_execution_time_ns; } dispatch->state.publish(root); - }).name("render_2d.paint.complete"); - paint_if.precede(paint_run, paint_done); + }); + paint_if.precede(paint_run); + paint_if.precede(paint_done); paint_run.precede(paint_extension); paint_extension.precede(paint_done); if (has_previous) previous.precede(paint_if); previous = paint_done; if (paint_node.cache_group_last) { Root* cache_owner = paint_node.cache_owner; - auto composite = taskflow.emplace([this, object, cache_owner] { + auto composite = taskflow.add(task_prefix + ".cache_composite", [this, object, cache_owner] { const auto group = std::find_if(cache_groups.begin(), cache_groups.end(), [cache_owner](const Cache_Group& value) { return value.owner == cache_owner; }); const auto owner_node = std::find_if(paint_order.begin(), paint_order.end(), [cache_owner](const Paint_Node& value) { return value.object == cache_owner; }); if (group == cache_groups.end() || owner_node == paint_order.end()) return; @@ -204,7 +205,7 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access if (!frame_target) throw std::logic_error("2D paint graph has no external frame target"); frame_target->composite(*owner_node->private_data->valid_cache); } - }).name("render_2d.paint.cache.composite"); + }); paint_done.precede(composite); previous = composite; } @@ -295,7 +296,13 @@ void Render_Scene_2D::Private::process(Object* object, Callback&& callback) Pen{.style = Line_Style::none}, Brush{state.background, Brush_Style::solid}); } prepare_paint_targets(object, frame, state.viewport); - if (paint_taskflow && !paint_taskflow->empty()) aethera::detail::run_taskflow(*paint_taskflow); + if (paint_taskflow && !paint_taskflow->empty()) { + if (frame_object->taskflow_trace_requested()) + aethera::detail::run_taskflow(*paint_taskflow, *frame_object, + "render_2d.paint"); + else + aethera::detail::run_taskflow(*paint_taskflow); + } frame_object->mark(Frame_Trace_Marker::paint_finished); }); std::invoke(std::forward(callback)); @@ -322,6 +329,14 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { } } admission{*this}; frame->mark(Frame_Trace_Marker::scene_render_requested); + const bool trace_started = aethera::detail::begin_taskflow_trace(*frame); + struct Trace_Scope { + Frame_2D* frame{}; /* 异常路径仍需关闭本帧 Observer 写入窗口。 */ + bool active{}; /* 只有成功取得全局捕获槽位才负责关闭。 */ + ~Trace_Scope() { + if (active) aethera::detail::finish_taskflow_trace(*frame); + } + } trace_scope{frame, trace_started}; struct Active_Frame_Scope { Frame_2D*& target; /* 最终 Private 的同步 process 帧槽位。 */ Frame_2D* previous{}; /* 嵌套调用前的帧;析构时恢复。 */ @@ -333,7 +348,17 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { completed = true; frame->mark(Frame_Trace_Marker::scene_render_finished); /* 帧准入直到 completion 图和最终回调都结束才释放,因此图运行期间外部帧不会被下一帧复用。 */ - if (!completion_graph.empty()) aethera::detail::run_taskflow(completion_graph); + if (!completion_graph.empty()) { + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow(completion_graph, *frame, + "render_2d.completion"); + else + aethera::detail::run_taskflow(completion_graph); + } + if (trace_started) { + aethera::detail::finish_taskflow_trace(*frame); + trace_scope.active = false; + } frame->mark(Frame_Trace_Marker::callback_started); callback(frame); frame->mark(Frame_Trace_Marker::callback_finished); diff --git a/render_2D/tests/Axis_Test.cpp b/render_2D/tests/Axis_Test.cpp index 04ad208..bb81991 100644 --- a/render_2D/tests/Axis_Test.cpp +++ b/render_2D/tests/Axis_Test.cpp @@ -95,12 +95,12 @@ TEST(axis_render, scene_prepares_and_paints_uncached_axis_into_frame) { }).has_value())); int paint_extension_calls{}; int completion_calls{}; - axis->paint_taskflow().emplace( - [&] { ++paint_extension_calls; }).name("test.axis.paint.extension"); - scene->completion_taskflow().emplace([&] { + axis->paint_taskflow().add( + "test.axis.paint.extension", [&] { ++paint_extension_calls; }); + scene->completion_taskflow().add("test.scene_2d.completion", [&] { EXPECT_EQ(paint_extension_calls, 1); ++completion_calls; - }).name("test.scene_2d.completion"); + }); const auto frame = render_frame(scene.get()); EXPECT_EQ(paint_extension_calls, 1); EXPECT_EQ(completion_calls, 1); diff --git a/render_3D/render_3D/detail/Async_Render_Backend.cpp b/render_3D/render_3D/detail/Async_Render_Backend.cpp index 723db1c..372ed64 100644 --- a/render_3D/render_3D/detail/Async_Render_Backend.cpp +++ b/render_3D/render_3D/detail/Async_Render_Backend.cpp @@ -185,7 +185,8 @@ void Async_Render_Backend::Implementation::enqueue(Command command) { } if (!schedule) return; auto self = shared_from_this(); - aethera::schedule_task([self] { self->drain(); }); + aethera::schedule_task("render_3d.backend.drain", + [self] { self->drain(); }); } void Async_Render_Backend::Implementation::drain() noexcept { for (;;) { diff --git a/render_3D/render_3D/scene/Render_Scene_3D.cpp b/render_3D/render_3D/scene/Render_Scene_3D.cpp index 433f694..473cd72 100644 --- a/render_3D/render_3D/scene/Render_Scene_3D.cpp +++ b/render_3D/render_3D/scene/Render_Scene_3D.cpp @@ -10,7 +10,7 @@ bool Render_Scene_3D::Prop::operator==(const Prop&) const = default; bool Render_Scene_3D::State::operator==(const State&) const = default; Render_Scene_3D::Render_Result Render_Scene_3D::render(Frame_3D* frame) { return static_cast(*d).dispatch->render(this, frame); } void Render_Scene_3D::set_frame_callback(Frame_Callback callback) { static_cast(*d).dispatch->set_frame_callback(this, std::move(callback)); } -tf::Taskflow& Render_Scene_3D::completion_taskflow() { +Task_Graph& Render_Scene_3D::completion_taskflow() { return static_cast(*d).completion_graph; } void Render_Scene_3D::reset_frame_statistics() { diff --git a/render_3D/render_3D/scene/Render_Scene_3D.hpp b/render_3D/render_3D/scene/Render_Scene_3D.hpp index 81144aa..e4a3b78 100644 --- a/render_3D/render_3D/scene/Render_Scene_3D.hpp +++ b/render_3D/render_3D/scene/Render_Scene_3D.hpp @@ -70,7 +70,7 @@ struct Render_Scene_3D : Defmark(Frame_Trace_Marker::callback_started); if (callback) callback(frame); frame->mark(Frame_Trace_Marker::callback_finished); @@ -211,7 +212,11 @@ void Render_Scene_3D::Private::initialize_backend( * GPU 完成线程只触发 Taskflow topology,不等待编码/发送节点;最终回调在 topology 完成后执行。 * 帧准入在 finish() 之前始终有效,因此 completion 图运行期间该外部 Frame 不会被下一帧复用。 */ - aethera::detail::run_taskflow(completion_graph, std::move(finish)); + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow(completion_graph, *frame, + "render_3d.completion", std::move(finish)); + else + aethera::detail::run_taskflow(completion_graph, std::move(finish)); }); paint_context = std::make_shared>( detail::Scene_Paint_Context{ @@ -290,14 +295,27 @@ void Render_Scene_3D::Private::process(Object* object, Callback&& callback) requ if (!state.paint_executed) return; const auto paint_started = std::chrono::steady_clock::now(); if (dispatch->paint.builder) { - if (!data->paint_graph) data->paint_graph = std::make_unique(); + if (!data->paint_graph) data->paint_graph = std::make_unique( + std::string(dispatch->business_name) + ".paint"); const bool rebuild = dispatch->paint.rebuild_predicate(root); if (!data->paint_graph_built || rebuild) { *data->paint_graph = dispatch->paint.builder(root); data->paint_graph_built = true; state.paint_graph_rebuilt = true; } - state.paint_task_count = data->paint_graph->num_tasks(); - if (!data->paint_graph->empty()) aethera::detail::run_taskflow(*data->paint_graph); + state.paint_task_count = data->paint_graph->size(); + if (!data->paint_graph->empty()) { + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow( + *data->paint_graph, *frame, + std::string(dispatch->business_name) + ".paint"); + else + aethera::detail::run_taskflow(*data->paint_graph); + } } else { state.paint_task_count = dispatch->paint.run ? 1 : 0; if (dispatch->paint.run) dispatch->paint.run(root); } if (!data->paint_extension.empty()) - aethera::detail::run_taskflow(data->paint_extension); + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow( + data->paint_extension, *frame, + std::string(dispatch->business_name) + ".paint.extension"); + else + aethera::detail::run_taskflow(data->paint_extension); root->template take_dirty(); state.paint_execution_time_ns = static_cast( std::chrono::duration_cast( @@ -336,6 +354,14 @@ Render_Scene_3D::Render_Result Render_Scene_3D::Private::render(Object* object, } } admission{*this, submitted}; frame->mark(Frame_Trace_Marker::scene_render_requested); + const bool trace_started = aethera::detail::begin_taskflow_trace(*frame); + struct Trace_Scope { + Frame_3D* frame{}; /* 提交前异常或拒绝时关闭 Observer 写入窗口。 */ + bool active{}; /* 提交成功后所有权转交异步完成回调。 */ + ~Trace_Scope() { + if (active) aethera::detail::finish_taskflow_trace(*frame); + } + } trace_scope{frame, trace_started}; struct Active_Frame_Scope { Frame_3D*& target; /* 最终 Private 的同步 process 帧槽位。 */ Frame_3D* previous{}; /* 嵌套调用前的帧;析构时恢复。 */ @@ -343,7 +369,10 @@ Render_Scene_3D::Render_Result Render_Scene_3D::Private::render(Object* object, } active_frame_scope{active_frame, active_frame}; active_frame = frame; object->process([&] { submitted = true; frame->mark(Frame_Trace_Marker::scene_render_finished); }); - if (submitted) return Render_Result::submitted; + if (submitted) { + trace_scope.active = false; + return Render_Result::submitted; + } const auto& prop = object->template read_prop(); if (!prop.view_active) return Render_Result::view_inactive; if (prop.viewport.empty()) return Render_Result::empty_viewport; diff --git a/render_3D/render_3D/visual/Basic_Visual.ipp b/render_3D/render_3D/visual/Basic_Visual.ipp index 487b062..acb8f95 100644 --- a/render_3D/render_3D/visual/Basic_Visual.ipp +++ b/render_3D/render_3D/visual/Basic_Visual.ipp @@ -19,7 +19,7 @@ struct Basic_Visual::Private : Basic_Visual::Prev_Private { void prepare_data(Object* object); /* CRTP 子图能力:提交本 Visual 最近一次 Prepare 生成的唯一产物。 */ template - [[nodiscard]] tf::Taskflow build_paint_graph(Object* object, const Prop& prop); + [[nodiscard]] Task_Graph build_paint_graph(Object* object, const Prop& prop); /* CRTP 覆盖:绑定 Scene 后每帧执行 Paint;dirty 只表示 Prepared 数据是否更新,不控制帧生成。 */ template [[nodiscard]] bool should_paint(Object* object, const State& state, bool dirty); @@ -67,13 +67,13 @@ void Basic_Visual::Private::prepare_data(Object* object) { } template template -tf::Taskflow Basic_Visual::Private::build_paint_graph(Object* object, const Prop&) { - tf::Taskflow graph; - graph.emplace([this, object] { +Task_Graph Basic_Visual::Private::build_paint_graph(Object* object, const Prop&) { + Task_Graph graph{"render_3d.visual.paint"}; + graph.add("publish", [this, object] { if (auto context = paint_target.context.lock(); context && paint_target.enqueue) { paint_target.enqueue(context.get(), object, prepared); } - }).name("render_3d.paint.publish"); + }); return graph; } template diff --git a/render_3D/tests/Visual_Tests.cpp b/render_3D/tests/Visual_Tests.cpp index 58f314a..48d04db 100644 --- a/render_3D/tests/Visual_Tests.cpp +++ b/render_3D/tests/Visual_Tests.cpp @@ -68,9 +68,9 @@ TEST(Render_3D_Scene, Paint_Submits_Without_Waiting_For_Gpu) { Frame_3D exchanged_unchanged_frame{Frame_Identity{3, 0}}; Frame_3D* completed_frame{}; std::atomic_size_t completion_calls{}; - scene->completion_taskflow().emplace([&] { + scene->completion_taskflow().add("test.scene_3d.completion", [&] { completion_calls.fetch_add(1, std::memory_order_release); - }).name("test.scene_3d.completion"); + }); scene->set_frame_callback([&](Frame_3D* value) { EXPECT_GT(completion_calls.load(std::memory_order_acquire), 0u); completed_frame = value; diff --git a/web_server/src/Gallery_Video_Stream.cpp b/web_server/src/Gallery_Video_Stream.cpp index ed6fd84..cbd45f4 100644 --- a/web_server/src/Gallery_Video_Stream.cpp +++ b/web_server/src/Gallery_Video_Stream.cpp @@ -200,7 +200,7 @@ struct Gallery_Video_Stream::Private { } } if (!schedule) return; - aethera::schedule_task([lifetime] { + aethera::schedule_task("web.gallery.encode", [lifetime] { if (const auto owner = lifetime.lock()) owner->d->encode_latest(lifetime); }); @@ -361,7 +361,7 @@ struct Gallery_Video_Stream::Private { } if (again) { try { - aethera::schedule_task([lifetime] { + aethera::schedule_task("web.gallery.encode.continue", [lifetime] { if (const auto owner = lifetime.lock()) owner->d->encode_latest(lifetime); }); @@ -397,6 +397,7 @@ void Gallery_Video_Stream::bind_plots() { [weak](std::exception_ptr failure) { if (const auto owner = weak.lock()) { aethera::schedule_task( + "web.gallery.failure", [weak, failure = std::move(failure)]() mutable { if (const auto stream = weak.lock()) stream->d->fail(std::move(failure)); diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index 5a3691c..0cefff5 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -202,6 +202,54 @@ void append_event_statistics_json(nlohmann::json& output, } } +nlohmann::json taskflow_trace_json(const Taskflow_Frame_Trace& trace) { + nlohmann::json graphs = nlohmann::json::array(); + std::unordered_map node_ids; + for (const auto& graph : trace.graphs) { + nlohmann::json nodes = nlohmann::json::array(); + for (const auto& node : graph.nodes) { + node_ids.emplace(node.native_id, node.node_id); + nlohmann::json predecessors = nlohmann::json::array(); + for (const auto native_id : node.predecessors) + predecessors.push_back(std::to_string(native_id)); + nlohmann::json successors = nlohmann::json::array(); + for (const auto native_id : node.successors) + successors.push_back(std::to_string(native_id)); + nodes.push_back({ + {"native_id", std::to_string(node.native_id)}, {"id", node.node_id}, + {"parent_id", node.parent_node_id}, {"name", node.name}, + {"type", node.type}, {"predecessors", std::move(predecessors)}, + {"successors", std::move(successors)}}); + } + graphs.push_back({ + {"stage", graph.stage}, {"name", graph.taskflow_name}, + {"submitted_ms", graph.submitted_ms}, + {"finished_ms", graph.finished_ms}, + {"completed", graph.completed}, {"nodes", std::move(nodes)}}); + } + nlohmann::json executions = nlohmann::json::array(); + for (const auto& task : trace.tasks) { + const auto found = node_ids.find(task.native_id); + executions.push_back({ + {"native_id", std::to_string(task.native_id)}, + {"node_id", found == node_ids.end() ? std::string{} : found->second}, + {"worker_id", task.worker_id}, + {"worker_queue_size", task.worker_queue_size}, + {"worker_queue_capacity", task.worker_queue_capacity}, + {"ready_ms", task.ready_ms}, {"started_ms", task.started_ms}, + {"finished_ms", task.finished_ms}, + {"duration_ms", task.duration_ms}, + {"queue_wait_ms", task.queue_wait_ms}}); + } + return { + {"sequence", trace.identity.sequence}, + {"correlation_id", trace.identity.correlation_id}, + {"created_time_unix_ns", trace.created_time_unix_ns}, + {"worker_count", trace.worker_count}, + {"graphs", std::move(graphs)}, + {"executions", std::move(executions)}}; +} + template void dispatch_plot_input(Scene_Object& scene, const Plot_Input_Event& input) { @@ -299,6 +347,17 @@ struct Plot::Private { std::mutex render_preparation_mutex{}; /* 仅串行本图 CPU prepare;Scene 管理帧在途。 */ double last_clock_render_time_ms{-std::numeric_limits::infinity()}; std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()}; + std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */ + std::atomic_uint64_t coalesced_tick_count{}; /* 尚未消费时被更新 tick 替换的旧 tick 总数。 */ + std::atomic_uint64_t policy_skip_count{}; /* 帧策略拒绝的 tick 总数。 */ + std::atomic_uint64_t preparation_busy_count{}; /* 本 Plot Prepare 已在执行而跳过的提交次数。 */ + std::atomic_uint64_t frame_slot_busy_count{}; /* 三个物理帧槽均被占用的提交次数。 */ + std::atomic_uint64_t scene_rejection_count{}; /* Scene 单帧准入拒绝的提交次数。 */ + std::atomic_uint64_t submitted_frame_count{}; /* 成功提交给 Scene 的帧总数。 */ + std::atomic_size_t taskflow_trace_remaining{}; /* 尚待标记的实际渲染帧数。 */ + mutable std::mutex taskflow_trace_mutex{}; /* 只保护低频请求结果的交换。 */ + std::size_t taskflow_trace_requested_count{}; /* 当前批次请求总帧数。 */ + std::vector taskflow_traces{}; /* 已完成帧直接发布的 DAG 与 Observer 结果。 */ template Private(std::unique_ptr value_scene, @@ -319,6 +378,7 @@ struct Plot::Private { void clock_tick(const Plot_Render_Tick& tick); void render_frame(Plot_Render_Tick tick); void queue_completed_frame(Render_Frame* frame); + [[nodiscard]] bool mark_taskflow_trace(Render_Frame& frame); void fail(std::exception_ptr failure) noexcept; }; @@ -403,7 +463,7 @@ void Plot::Private::consume_tick(std::weak_ptr lifetime) { if (!schedule_again) tick_task_scheduled = false; } if (schedule_again) { - aethera::schedule_task([lifetime] { + aethera::schedule_task("web.plot.tick.consume", [lifetime] { const auto plot = lifetime.lock(); if (!plot) return; try { plot->d->consume_tick(lifetime); } @@ -415,20 +475,42 @@ void Plot::Private::consume_tick(std::weak_ptr lifetime) { void Plot::Private::clock_tick(const Plot_Render_Tick& tick) { if (terminal_failure.load(std::memory_order_acquire)) return; const auto pacing = frame_policy.snapshot(); - if (!pacing.render_enabled || pacing.mode == Frame_Pacing_Mode::manual) return; + if (!pacing.render_enabled || pacing.mode == Frame_Pacing_Mode::manual) { + policy_skip_count.fetch_add(1, std::memory_order_relaxed); + return; + } if (pacing.mode == Frame_Pacing_Mode::fixed_rate) { if (tick.time_milliseconds < last_clock_render_time_ms) last_clock_render_time_ms = -std::numeric_limits::infinity(); const double interval = 1'000.0 / pacing.fixed_rate_fps; - if (tick.time_milliseconds - last_clock_render_time_ms + 0.01 < interval) return; + if (tick.time_milliseconds - last_clock_render_time_ms + 0.01 < interval) { + policy_skip_count.fetch_add(1, std::memory_order_relaxed); + return; + } } last_clock_render_time_ms = tick.time_milliseconds; render_frame(tick); } +bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) { + auto remaining = taskflow_trace_remaining.load(std::memory_order_acquire); + while (remaining != 0) { + if (taskflow_trace_remaining.compare_exchange_weak( + remaining, remaining - 1, std::memory_order_acq_rel, + std::memory_order_acquire)) { + frame.request_taskflow_trace(); + return true; + } + } + return false; +} + void Plot::Private::render_frame(Plot_Render_Tick tick) { std::unique_lock preparation_lock(render_preparation_mutex, std::try_to_lock); - if (!preparation_lock.owns_lock()) return; + if (!preparation_lock.owns_lock()) { + preparation_busy_count.fetch_add(1, std::memory_order_relaxed); + return; + } if (terminal_failure.load(std::memory_order_acquire)) return; const auto streams = stream_snapshot(); const auto pacing = frame_policy.snapshot(); @@ -437,7 +519,10 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { std::lock_guard lock(frame_mutex); if (std::ranges::none_of(frame_slots, [](const Managed_Frame& slot) { return slot.state == Frame_State::available; - })) return; + })) { + frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed); + return; + } } tick.width = streams.width; tick.height = streams.height; @@ -446,7 +531,22 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { * 进入 Scene::render 后只剩已经准备好的 Visual 批次与轻量提交; * 共享 Render Domain 不承担业务数据生成。 */ + const auto update_started = std::chrono::steady_clock::now(); view->update(tick); + const auto update_elapsed = std::chrono::steady_clock::now() - update_started; + const auto tick_queue_elapsed = tick.issued_at.time_since_epoch().count() == 0 + ? std::chrono::steady_clock::duration::zero() + : update_started - tick.issued_at; + const auto record_plot_measurements = [&](Render_Frame& frame) { + const auto nanoseconds = [](std::chrono::steady_clock::duration duration) { + return static_cast(std::max(0, + std::chrono::duration_cast(duration).count())); + }; + frame.record(Frame_Trace_Measurement::plot_tick_queue_ns, + nanoseconds(tick_queue_elapsed)); + frame.record(Frame_Trace_Measurement::plot_update_ns, + nanoseconds(update_elapsed)); + }; const std::uint64_t sequence = next_frame_sequence++; const Frame_Identity identity{sequence, tick.sequence == 0 ? sequence : tick.sequence}; @@ -474,6 +574,11 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { if (slot.state == Frame_State::in_flight) slot.state = Frame_State::available; }; + bool taskflow_trace_claimed{}; + const auto restore_taskflow_trace_claim = [this, &taskflow_trace_claimed] { + if (!std::exchange(taskflow_trace_claimed, false)) return; + taskflow_trace_remaining.fetch_add(1, std::memory_order_release); + }; try { if (auto* scene_2d = std::get_if>(&scene)) { @@ -481,26 +586,44 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { output.begin(identity, pacing.video_enabled ? render_2d::Pixel_Format::rgba8 : Frame_2D::native_pixel_format); + taskflow_trace_claimed = mark_taskflow_trace(output); + record_plot_measurements(output); (*scene_2d)->set<&Render_Scene_2D::Prop::viewport>( Size{static_cast(tick.width), static_cast(tick.height)}); const auto result = (*scene_2d)->render(&output); - if (!result) rollback_unsubmitted(); + if (!result) { + scene_rejection_count.fetch_add(1, std::memory_order_relaxed); + rollback_unsubmitted(); + restore_taskflow_trace_claim(); + } else { + taskflow_trace_claimed = false; + submitted_frame_count.fetch_add(1, std::memory_order_relaxed); + } return; } auto& output = *std::get>(managed->frame); output.begin(identity, pacing.video_enabled ? Frame_3D_Output::pixels : Frame_3D_Output::diagnostics, Frame_3D::native_pixel_format); + taskflow_trace_claimed = mark_taskflow_trace(output); + record_plot_measurements(output); auto& scene_3d = std::get>(scene); scene_3d->set<&Render_Scene_3D::Prop::viewport>(Extent{tick.width, tick.height}); const auto result = scene_3d->render(&output); - if (result == Render_Scene_3D::Render_Result::submitted) return; + if (result == Render_Scene_3D::Render_Result::submitted) { + taskflow_trace_claimed = false; + submitted_frame_count.fetch_add(1, std::memory_order_relaxed); + return; + } + scene_rejection_count.fetch_add(1, std::memory_order_relaxed); rollback_unsubmitted(); + restore_taskflow_trace_claim(); if (result == Render_Scene_3D::Render_Result::backend_unavailable) throw std::runtime_error("3D render backend became unavailable before submission"); } catch (...) { rollback_unsubmitted(); + restore_taskflow_trace_claim(); throw; } } @@ -550,6 +673,11 @@ void Plot::Private::queue_completed_frame(Render_Frame* frame) { }; try { + if (frame->taskflow_trace_requested()) { + auto trace = frame->taskflow_trace(); + std::lock_guard lock(taskflow_trace_mutex); + taskflow_traces.push_back(std::move(trace)); + } const auto pacing = frame_policy.snapshot(); const auto identity = frame->identity(); Frame_Identity rendered_identity = identity; @@ -589,7 +717,12 @@ void Plot::Private::queue_completed_frame(Render_Frame* frame) { const auto published = std::make_shared( Plot_Stream_Frame{{}, std::move(pixels)}); + const auto publish_started = std::chrono::steady_clock::now(); publish(std::move(published)); + frame->record(Frame_Trace_Measurement::plot_publish_ns, + static_cast(std::max(0, + std::chrono::duration_cast( + std::chrono::steady_clock::now() - publish_started).count()))); retire(); } catch (...) { @@ -680,9 +813,12 @@ void Plot::configure_stream(Stream_Id stream, std::uint32_t width, void Plot::schedule_render(Plot_Render_Tick tick) { ensure_started(); if (d->terminal_failure.load(std::memory_order_acquire)) return; + d->received_tick_count.fetch_add(1, std::memory_order_relaxed); bool schedule{}; { std::lock_guard lock(d->tick_mutex); + if (d->pending_tick) + d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed); d->pending_tick = tick; if (!d->tick_task_scheduled) { d->tick_task_scheduled = true; @@ -691,7 +827,7 @@ void Plot::schedule_render(Plot_Render_Tick tick) { } if (!schedule) return; const auto weak = weak_from_this(); - aethera::schedule_task([weak] { + aethera::schedule_task("web.plot.tick.consume", [weak] { const auto owner = weak.lock(); if (!owner) return; try { @@ -706,13 +842,14 @@ void Plot::schedule_render(Plot_Render_Tick tick) { void Plot::render_once() { ensure_started(); const auto weak = weak_from_this(); - aethera::schedule_task([weak] { + aethera::schedule_task("web.plot.render.once", [weak] { const auto owner = weak.lock(); if (!owner) return; try { const auto elapsed = std::chrono::steady_clock::now() - owner->d->clock_origin; owner->d->render_frame(Plot_Render_Tick{ + std::chrono::steady_clock::now(), 0, std::chrono::duration(elapsed).count()}); } catch (...) { @@ -832,6 +969,14 @@ nlohmann::json Plot::diagnostics() const { {"fixed_rate_fps", pacing.fixed_rate_fps}, {"render_enabled", pacing.render_enabled}, {"video_enabled", pacing.video_enabled}}}, + {"plot_scheduler", { + {"received_ticks", d->received_tick_count.load(std::memory_order_relaxed)}, + {"coalesced_ticks", d->coalesced_tick_count.load(std::memory_order_relaxed)}, + {"policy_skips", d->policy_skip_count.load(std::memory_order_relaxed)}, + {"preparation_busy", d->preparation_busy_count.load(std::memory_order_relaxed)}, + {"frame_slot_busy", d->frame_slot_busy_count.load(std::memory_order_relaxed)}, + {"scene_rejections", d->scene_rejection_count.load(std::memory_order_relaxed)}, + {"submitted_frames", d->submitted_frame_count.load(std::memory_order_relaxed)}}}, {"frame_statistics", std::move(frame_statistics)}, {"input_statistics", std::move(input_statistics)}}; if (is_3d) { @@ -865,6 +1010,40 @@ nlohmann::json Plot::diagnostics() const { return output; } +void Plot::request_taskflow_trace(std::size_t frame_count) { + if (frame_count == 0 || frame_count > 120) + throw std::invalid_argument("Taskflow trace frame_count must be between 1 and 120"); + ensure_started(); + { + std::lock_guard lock(d->taskflow_trace_mutex); + if (d->taskflow_trace_requested_count != d->taskflow_traces.size()) + throw std::logic_error("A Taskflow frame trace request is already active"); + d->taskflow_traces.clear(); + d->taskflow_traces.reserve(frame_count); + d->taskflow_trace_requested_count = frame_count; + } + d->taskflow_trace_remaining.store(frame_count, std::memory_order_release); +} + +nlohmann::json Plot::taskflow_trace() const { + nlohmann::json frames = nlohmann::json::array(); + std::size_t requested{}; + { + std::lock_guard lock(d->taskflow_trace_mutex); + requested = d->taskflow_trace_requested_count; + for (const auto& trace : d->taskflow_traces) + frames.push_back(taskflow_trace_json(trace)); + } + const auto remaining = d->taskflow_trace_remaining.load( + std::memory_order_acquire); + return { + {"protocol", "aethera.taskflow.frames"}, {"version", 1}, + {"requested", requested}, {"remaining", remaining}, + {"captured", frames.size()}, + {"complete", requested != 0 && frames.size() == requested}, + {"frames", std::move(frames)}}; +} + void Plot::reset_diagnostics() { std::visit([](auto& scene) { scene->reset_frame_statistics(); }, d->scene); } diff --git a/web_server/src/Plot.hpp b/web_server/src/Plot.hpp index 47c5b93..cc3880d 100644 --- a/web_server/src/Plot.hpp +++ b/web_server/src/Plot.hpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -28,6 +29,7 @@ struct Plot_Input_Event { }; struct Plot_Render_Tick { + std::chrono::steady_clock::time_point issued_at{}; /* 页面帧时钟发布本 tick 的单调时刻;手动帧在提交时填写。 */ std::uint64_t sequence{}; /* 页面级帧时钟分配的关联序号。 */ double time_milliseconds{}; /* 页面级单调时间线,所有图共享同一个动画时刻。 */ std::uint32_t width{320}; /* 当前图在媒体图集中的固定像素宽度。 */ @@ -88,6 +90,9 @@ public: const nlohmann::json& value); [[nodiscard]] nlohmann::json generate_data(const nlohmann::json& input); [[nodiscard]] nlohmann::json diagnostics() const; + /* 清空旧捕获并请求接下来实际完成的 frame_count 帧 Task DAG。 */ + void request_taskflow_trace(std::size_t frame_count); + [[nodiscard]] nlohmann::json taskflow_trace() const; void reset_diagnostics(); private: diff --git a/web_server/src/Web_Server.cpp b/web_server/src/Web_Server.cpp index dfe850f..032a533 100644 --- a/web_server/src/Web_Server.cpp +++ b/web_server/src/Web_Server.cpp @@ -38,6 +38,64 @@ std::shared_ptr find_plot(const Plot_Map& plots, std::string_view id) { const auto found = plots.find(std::string(id)); return found == plots.end() ? nullptr : found->second; } + +nlohmann::json taskflow_runtime_json() { + const auto state = aethera::task_runtime_state(); + nlohmann::json workers = nlohmann::json::array(); + for (const auto& worker : state.workers) + workers.push_back({ + {"id", worker.id}, {"task_count", worker.task_count}, + {"current_queue_size", worker.current_queue_size}, + {"current_queue_capacity", worker.current_queue_capacity}, + {"peak_queue_size", worker.peak_observed_queue_size}, + {"max_queue_capacity", worker.max_observed_queue_capacity}, + {"active_task", {{"native_id", std::to_string(worker.active_task_hash)}, + {"type", worker.active_task_type}, + {"time_ns", worker.active_task_time_ns}}}, + {"task_time_ns", worker.task_time_ns}, + {"busy_time_ns", worker.busy_time_ns}, + {"idle_time_ns", worker.idle_time_ns}, + {"min_task_time_ns", worker.min_task_time_ns}, + {"max_task_time_ns", worker.max_task_time_ns}, + {"utilization", worker.utilization}}); + nlohmann::json task_types = nlohmann::json::array(); + for (const auto& type : state.task_types) + task_types.push_back({ + {"name", type.name}, {"count", type.count}, + {"total_time_ns", type.total_time_ns}, + {"min_time_ns", type.min_time_ns}, + {"max_time_ns", type.max_time_ns}}); + return { + {"protocol", "aethera.taskflow.runtime"}, {"version", 1}, + {"worker_count", state.worker_count}, + {"active_topologies", state.active_topology_count}, + {"active_taskflows", state.active_taskflow_count}, + {"peak_active_taskflows", state.peak_active_taskflow_count}, + {"completed_taskflows", state.completed_taskflow_count}, + {"failed_taskflows", state.failed_taskflow_count}, + {"active_tasks", state.active_task_count}, + {"peak_active_tasks", state.peak_active_task_count}, + {"active_workers", state.active_worker_count}, + {"peak_active_workers", state.peak_active_worker_count}, + {"observed_tasks", state.observed_task_count}, + {"named_tasks", state.named_task_count}, + {"peak_worker_queue_size", state.peak_observed_worker_queue_size}, + {"max_worker_queue_capacity", state.max_observed_worker_queue_capacity}, + {"max_predecessors", state.max_predecessors}, + {"max_successors", state.max_successors}, + {"max_strong_dependencies", state.max_strong_dependencies}, + {"max_weak_dependencies", state.max_weak_dependencies}, + {"longest_task", {{"native_id", std::to_string(state.longest_task_hash)}, + {"name", state.longest_task_name}, + {"type", state.longest_task_type}, + {"time_ns", state.longest_task_time_ns}}}, + {"total_task_time_ns", state.total_task_time_ns}, + {"worker_busy_time_ns", state.worker_busy_time_ns}, + {"observed_wall_time_ns", state.observed_wall_time_ns}, + {"worker_utilization", state.worker_utilization}, + {"task_types", std::move(task_types)}, + {"workers", std::move(workers)}}; +} } int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) { @@ -131,7 +189,8 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) {"dimension", id.starts_with("datoviz_") ? "3D" : "2D"}, {"websocket", "/ws/plot/" + id}, {"media", plot_media->at(id)}, {"schema", "/plot/" + id + "/schema"}, - {"diagnostics", "/plot/" + id + "/diagnostics"}}); + {"diagnostics", "/plot/" + id + "/diagnostics"}, + {"taskflow", "/plot/" + id + "/taskflow"}}); } callback(json_response(std::move(result))); }, {drogon::Get}); @@ -153,6 +212,46 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) callback(json_response(plot->diagnostics())); }, {drogon::Get, drogon::Delete}); + app.registerHandler("/taskflow/diagnostics", []( + const drogon::HttpRequestPtr&, + std::function&& callback) { + callback(json_response(taskflow_runtime_json())); + }, {drogon::Get}); + + app.registerHandler("/plot/{1}/taskflow", [plots]( + const drogon::HttpRequestPtr& request, + std::function&& callback, + std::string plot_id) { + auto plot = find_plot(*plots, plot_id); + if (!plot) { + callback(error_response(drogon::k404NotFound, "unknown plot")); + return; + } + try { + if (request->method() == drogon::Post) { + const auto input = nlohmann::json::parse(request->body()); + if (!input.is_object() || !input.contains("frame_count") || + !input["frame_count"].is_number_unsigned()) { + callback(error_response(drogon::k400BadRequest, + "Taskflow trace requires unsigned frame_count")); + return; + } + plot->request_taskflow_trace(input["frame_count"].get()); + } + callback(json_response(plot->taskflow_trace())); + } + catch (const nlohmann::json::exception&) { + callback(error_response(drogon::k400BadRequest, + "invalid Taskflow trace request")); + } + catch (const std::invalid_argument& failure) { + callback(error_response(drogon::k400BadRequest, failure.what())); + } + catch (const std::logic_error& failure) { + callback(error_response(drogon::k409Conflict, failure.what())); + } + }, {drogon::Get, drogon::Post}); + app.registerHandler("/gallery/{1}/diagnostics", [gallery_streams]( const drogon::HttpRequestPtr&, std::function&& callback, diff --git a/web_server/src/detail/Gallery_Frame_Clock.cpp b/web_server/src/detail/Gallery_Frame_Clock.cpp index ef25f67..9ccab16 100644 --- a/web_server/src/detail/Gallery_Frame_Clock.cpp +++ b/web_server/src/detail/Gallery_Frame_Clock.cpp @@ -51,6 +51,7 @@ struct Gallery_Frame_Clock::Private { try { ++sequence; tick_handler(Plot_Render_Tick{ + std::chrono::steady_clock::now(), sequence, std::chrono::duration(deadline - origin).count()}); deadline += interval; diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index 7f3bbbb..515fd97 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -2,6 +2,9 @@ import {memo, useCallback, useEffect, useMemo, useRef, useState} from "react"; import {I18nLabel, Layout, Model, type IJsonModel, type TabNode} from "flexlayout-react"; import {Responsive, useContainerWidth, type LayoutItem, type ResponsiveLayouts} from "react-grid-layout"; import ReconnectingWebSocket from "reconnecting-websocket"; +import ELK from "elkjs/lib/elk.bundled.js"; +import {Background, Controls, MarkerType, MiniMap, ReactFlow, + type Edge as Flow_Edge, type Node as Flow_Node} from "@xyflow/react"; import * as echarts from "echarts/core"; import {LineChart} from "echarts/charts"; import {DataZoomComponent, GridComponent, LegendComponent, TooltipComponent} from "echarts/components"; @@ -10,10 +13,11 @@ import type {EChartsType} from "echarts/core"; import "flexlayout-react/style/alpha_dark.css"; import "react-grid-layout/css/styles.css"; import "react-resizable/css/styles.css"; +import "@xyflow/react/dist/style.css"; echarts.use([LineChart, DataZoomComponent, GridComponent, LegendComponent, TooltipComponent, CanvasRenderer]); -type Plot = {id: string; title: string; category: string; description: string; dimension: "2D" | "3D"; websocket: string; media: string; schema: string; diagnostics: string}; +type Plot = {id: string; title: string; category: string; description: string; dimension: "2D" | "3D"; websocket: string; media: string; schema: string; diagnostics: string; taskflow: string}; type Option = {value: string; label: string}; type Editor = "boolean" | "integer" | "number" | "text" | "select" | "color" | "point2" | "size" | "rect" | "range" | "pen" | "brush" | "font" | "color-map" | "vector3" | "matrix4" | "marker-list" | "surface-marker-list" | "json"; type Color_Channel_Scale = "normalized" | "byte"; @@ -69,6 +73,28 @@ type Frame_Metrics = {sequence: number; generated_time_unix_ms: number; server_c pacing_mode: Frame_Pacing_Mode; fixed_rate_fps: number; delivery: Frame_Delivery; video_playback: Video_Playback_Metrics}; type Stage_Statistic = "average" | "variability" | "p95" | "p99"; type Stage_Unit = "value" | "percentage"; +type Taskflow_Node_Trace = {native_id: string; id: string; parent_id: string; name: string; type: string; + predecessors: string[]; successors: string[]}; +type Taskflow_Graph_Trace = {stage: string; name: string; submitted_ms: number; finished_ms: number; + completed: boolean; nodes: Taskflow_Node_Trace[]}; +type Taskflow_Execution_Trace = {native_id: string; node_id: string; worker_id: number; worker_queue_size: number; + worker_queue_capacity: number; ready_ms: number; started_ms: number; finished_ms: number; duration_ms: number; queue_wait_ms: number}; +type Taskflow_Frame_Trace = {sequence: number; correlation_id: number; created_time_unix_ns: number; worker_count: number; + graphs: Taskflow_Graph_Trace[]; executions: Taskflow_Execution_Trace[]}; +type Taskflow_Frame_Response = {protocol: "aethera.taskflow.frames"; version: 1; requested: number; remaining: number; + captured: number; complete: boolean; frames: Taskflow_Frame_Trace[]}; +type Taskflow_Worker_State = {id: number; task_count: number; current_queue_size: number; current_queue_capacity: number; + peak_queue_size: number; max_queue_capacity: number; active_task: {native_id: string; type: string; time_ns: number}; + task_time_ns: number; busy_time_ns: number; idle_time_ns: number; min_task_time_ns: number; max_task_time_ns: number; utilization: number}; +type Taskflow_Type_State = {name: string; count: number; total_time_ns: number; min_time_ns: number; max_time_ns: number}; +type Taskflow_Runtime_State = {protocol: "aethera.taskflow.runtime"; version: 1; worker_count: number; active_topologies: number; + active_taskflows: number; peak_active_taskflows: number; completed_taskflows: number; failed_taskflows: number; + active_tasks: number; peak_active_tasks: number; active_workers: number; peak_active_workers: number; observed_tasks: number; + named_tasks: number; peak_worker_queue_size: number; max_worker_queue_capacity: number; max_predecessors: number; + max_successors: number; max_strong_dependencies: number; max_weak_dependencies: number; + longest_task: {native_id: string; name: string; type: string; time_ns: number}; total_task_time_ns: number; + worker_busy_time_ns: number; observed_wall_time_ns: number; worker_utilization: number; + task_types: Taskflow_Type_State[]; workers: Taskflow_Worker_State[]}; const default_plot_execution_policy = (): Plot_Execution_Policy => ({visible: true}); @@ -1145,6 +1171,185 @@ function Frame_Statistics_Pane({plot, diagnostics, busy, on_refresh, on_reset}:
; } +const taskflow_layout = new ELK(); + +function milliseconds(value: number) { + if (!Number.isFinite(value)) return "--"; + return value < .01 ? `${(value * 1000).toFixed(1)} μs` : `${value.toFixed(3)} ms`; +} + +function nanoseconds(value: number) { + return milliseconds(value / 1_000_000); +} + +function Taskflow_Dag({graph, executions}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]}) { + const [nodes, set_nodes] = useState([]); + const [edges, set_edges] = useState([]); + useEffect(() => { + let cancelled = false; + const node_id = new Map(graph.nodes.map(node => [node.native_id, node.id])); + const execution = new Map(executions.map(value => [value.native_id, value])); + const flow_edges: Flow_Edge[] = []; + for (const node of graph.nodes) for (const successor of node.successors) { + const target = node_id.get(successor); + if (!target) continue; + flow_edges.push({id: `${node.id}->${target}`, source: node.id, target, + markerEnd: {type: MarkerType.ArrowClosed}, animated: false}); + } + void taskflow_layout.layout({ + id: "taskflow-frame", + layoutOptions: { + "elk.algorithm": "layered", "elk.direction": "RIGHT", + "elk.layered.spacing.nodeNodeBetweenLayers": "72", "elk.spacing.nodeNode": "34", + "elk.layered.nodePlacement.strategy": "NETWORK_SIMPLEX" + }, + children: graph.nodes.map(node => ({id: node.id, width: 230, height: 92})), + edges: flow_edges.map(edge => ({id: edge.id, sources: [edge.source], targets: [edge.target]})) + }).then(layout => { + if (cancelled) return; + set_nodes(graph.nodes.map(node => { + const position = layout.children?.find(item => item.id === node.id); + const sample = execution.get(node.native_id); + const wait = sample?.queue_wait_ms ?? 0; + const duration = sample?.duration_ms ?? 0; + return { + id: node.id, + position: {x: position?.x ?? 0, y: position?.y ?? 0}, + data: {label:
{node.name}{node.id} + {node.type} · W{sample?.worker_id ?? "--"} + 执行 {sample ? milliseconds(duration) : "未执行"} · 排队 {sample ? milliseconds(wait) : "--"}
}, + className: sample ? wait > duration && wait > .1 ? "taskflowNode taskflowNodeWaiting" : "taskflowNode taskflowNodeExecuted" : "taskflowNode" + }; + })); + set_edges(flow_edges); + }); + return () => { cancelled = true; }; + }, [graph, executions]); + return
+ +
; +} + +function Taskflow_Timeline({graph, executions}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]}) { + const native_ids = new Set(graph.nodes.map(node => node.native_id)); + const names = new Map(graph.nodes.map(node => [node.native_id, node.name])); + const rows = executions.filter(value => native_ids.has(value.native_id)).sort((left, right) => left.started_ms - right.started_ms); + const begin = Math.min(graph.submitted_ms, ...rows.map(row => row.ready_ms)); + const end = Math.max(graph.finished_ms, ...rows.map(row => row.finished_ms), begin + .001); + const duration = Math.max(.001, end - begin); + return
Worker 时间线横向位置按本帧真实 Observer 时间缩放;浅色段为估算就绪等待,亮色段为 on_entry → on_exit。
+ {rows.length ?
{rows.map((row, index) => { + const left = Math.max(0, (row.ready_ms - begin) / duration * 100); + const waiting = Math.max(.15, row.queue_wait_ms / duration * 100); + const running = Math.max(.25, row.duration_ms / duration * 100); + return
+ W{row.worker_id}{names.get(row.native_id) ?? row.node_id} +
+
+ {milliseconds(row.duration_ms)}{milliseconds(row.queue_wait_ms)} 排队 +
; + })}
:
本阶段没有 Observer 执行记录条件节点未命中或该图没有运行。
} +
; +} + +function Taskflow_Frame_Pane({plot}: {plot: Plot}) { + const [frame_count, set_frame_count] = useState(8); + const [response, set_response] = useState(null); + const [frame_index, set_frame_index] = useState(0); + const [graph_index, set_graph_index] = useState(0); + const [busy, set_busy] = useState(false); + const [error, set_error] = useState(""); + const load = useCallback(async () => { + const request = await fetch(plot.taskflow); + const value = await request.json() as Taskflow_Frame_Response & {error?: string}; + if (!request.ok) throw new Error(value.error ?? "读取 Taskflow 帧失败"); + set_response(value); + set_frame_index(current => Math.min(current, Math.max(0, value.frames.length - 1))); + }, [plot.taskflow]); + useEffect(() => { + set_response(null); set_frame_index(0); set_graph_index(0); set_error(""); + void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 帧失败")); + }, [load]); + useEffect(() => { + if (!response || response.requested === 0 || response.complete) return; + const timer = window.setInterval(() => void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 帧失败")), 400); + return () => window.clearInterval(timer); + }, [load, response?.requested, response?.complete]); + const capture = async () => { + set_busy(true); set_error(""); set_frame_index(0); set_graph_index(0); + try { + const request = await fetch(plot.taskflow, {method: "POST", headers: {"Content-Type": "application/json"}, body: JSON.stringify({frame_count})}); + const value = await request.json() as Taskflow_Frame_Response & {error?: string}; + if (!request.ok) throw new Error(value.error ?? "请求 Taskflow 帧失败"); + set_response(value); + } catch (failure) { + set_error(failure instanceof Error ? failure.message : "请求 Taskflow 帧失败"); + } finally { set_busy(false); } + }; + const frame = response?.frames[frame_index]; + const graph = frame?.graphs[graph_index]; + useEffect(() => set_graph_index(0), [frame?.sequence]); + return
void load()}/> +
+
+ + {response ? `${response.captured}/${response.requested} 帧${response.complete ? " · 已完成" : ` · 还需 ${response.remaining} 帧`}` : "按需捕获,未请求时 Observer 不写入逐帧数据"}
+ {error ?

{error}

: null} + {frame ? <> +
+ + Executor {frame.worker_count} workers · 本帧 {frame.executions.length} 次任务执行
+ {graph ? <>
+
业务阶段
{graph.stage}
DAG 节点
{graph.nodes.length}
+
Topology 总耗时
{milliseconds(graph.finished_ms - graph.submitted_ms)}
+
状态
{graph.completed ? "完成" : "未完成"}
+
:
该帧没有 Taskflow 阶段
} + :
等待逐帧 Taskflow 样本输入 N 后捕获后续实际渲染帧;每帧独立绑定其 DAG 元信息和 Observer 结果。
} +
; +} + +function Taskflow_Runtime_Pane() { + const [state, set_state] = useState(null); + const [error, set_error] = useState(""); + const load = useCallback(async () => { + const request = await fetch("/taskflow/diagnostics"); + const value = await request.json() as Taskflow_Runtime_State & {error?: string}; + if (!request.ok) throw new Error(value.error ?? "读取 Taskflow 总体状态失败"); + set_state(value); set_error(""); + }, []); + useEffect(() => { + void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 总体状态失败")); + const timer = window.setInterval(() => void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 总体状态失败")), 1000); + return () => window.clearInterval(timer); + }, [load]); + return
全局执行域

Taskflow 总体观测

每秒读取一次累计 Observer 状态;逐帧拓扑捕获在各图的“Taskflow 帧分析”页按需开启。

+
{error ?

{error}

: null}{state ? <> +
+
Worker
{state.worker_count}
总体占用率
{state.worker_utilization.toFixed(1)}%
+
活跃 Worker
{state.active_workers}/{state.worker_count}
活跃任务
{state.active_tasks}
+
活跃 Topology
{state.active_topologies}
活跃 Taskflow
{state.active_taskflows}
+
累计任务
{state.observed_tasks.toLocaleString("zh-CN")}
失败 Taskflow
{state.failed_taskflows}
+
队列峰值
{state.peak_worker_queue_size}
最长任务
{nanoseconds(state.longest_task.time_ns)}
+
+

Observer 的任务持续时间包含任务体内部的锁等待或阻塞;队列峰值与就绪等待能定位 Executor 拥塞,但 Taskflow 原生 Observer 不提供具体互斥量名称。具体帧请在逐帧 DAG 中核对并行关系、Worker 与长尾。

+
Worker 占用与队列累计值,不在 GET 时重新计算任务样本。
{state.workers.map(worker =>
+
Worker {worker.id}{worker.utilization.toFixed(1)}%
+
+
任务
{worker.task_count}
队列 当前/峰值
{worker.current_queue_size}/{worker.peak_queue_size}
+
最长
{nanoseconds(worker.max_task_time_ns)}
活跃持续
{worker.active_task.time_ns ? nanoseconds(worker.active_task.time_ns) : "空闲"}
+ {worker.active_task.time_ns ? {worker.active_task.type} · {worker.active_task.native_id} : null} +
)}
+
Taskflow 原生任务类型按 Observer TaskType 累计执行次数与耗时。
+ {state.task_types.map(type => + + )}
类型次数累计最短最长
{type.name}{type.count.toLocaleString("zh-CN")}{nanoseconds(type.total_time_ns)}{nanoseconds(type.min_time_ns)}{nanoseconds(type.max_time_ns)}
+ :
正在读取 Taskflow 总体状态
}
; +} + const Plot_Card = memo(function Plot_Card({plot, selected, policy, gallery, on_policy, on_select}: { plot: Plot; selected: boolean; policy: Plot_Execution_Policy; gallery: Gallery_Video_State; on_policy: (plot_id: string, patch: Partial) => void; on_select: (plot: Plot) => void; @@ -1315,7 +1520,7 @@ function Gallery_Grid({plots, selected, policies, layout_scope, galleries, on_po : null}; } -const workspace_layout_key = "aethera-flexlayout-v3"; +const workspace_layout_key = "aethera-flexlayout-v4"; const layout_labels: Record = { [I18nLabel.Close_Tab]: "关闭标签", [I18nLabel.Pinned_Tab]: "已固定", @@ -1361,7 +1566,9 @@ const default_workspace_layout: IJsonModel = { {type: "tab", id: "state-tab", name: "运行状态", component: "state", enableClose: false, enableScrollbars: false, minWidth: 320, minHeight: 240}, {type: "tab", id: "frame-policy-tab", name: "采样与帧策略", component: "frame-policy", enableClose: false, enableScrollbars: false, minWidth: 320, minHeight: 240}, {type: "tab", id: "data-generation-tab", name: "原始数据生成", component: "data-generation", enableClose: false, enableScrollbars: false, minWidth: 320, minHeight: 240}, - {type: "tab", id: "frame-statistics-tab", name: "帧流水线统计", component: "frame-statistics", enableClose: false, enableScrollbars: false, minWidth: 360, minHeight: 260} + {type: "tab", id: "frame-statistics-tab", name: "帧流水线统计", component: "frame-statistics", enableClose: false, enableScrollbars: false, minWidth: 360, minHeight: 260}, + {type: "tab", id: "taskflow-frame-tab", name: "Taskflow 帧分析", component: "taskflow-frame", enableClose: false, enableScrollbars: false, minWidth: 480, minHeight: 320}, + {type: "tab", id: "taskflow-runtime-tab", name: "Taskflow 总体", component: "taskflow-runtime", enableClose: false, enableScrollbars: false, minWidth: 480, minHeight: 320} ]} ]} }; @@ -1456,6 +1663,7 @@ export function App() { ; const factory = (node: TabNode) => { if (node.getComponent() === "gallery") return gallery; + if (node.getComponent() === "taskflow-runtime") return ; if (!selected) return
请选择一个图形组件。
; if (node.getComponent() === "properties") return ; if (node.getComponent() === "state") return ; @@ -1469,6 +1677,7 @@ export function App() { window.dispatchEvent(new CustomEvent("aethera-manual-frame", {detail: {plot_id: selected.id}})); }}/>; if (node.getComponent() === "frame-statistics") return ; + if (node.getComponent() === "taskflow-frame") return ; return
未知工作区面板。
; }; return
layout_labels[label]} diff --git a/webapp_gallery/src/styles.css b/webapp_gallery/src/styles.css index 26dd3fd..e731de4 100644 --- a/webapp_gallery/src/styles.css +++ b/webapp_gallery/src/styles.css @@ -182,6 +182,62 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .diagnosticEmpty { display: grid; min-height: 180px; place-content: center; gap: 8px; padding: 24px; color: #71839e; border: 1px dashed #29435e; border-radius: 10px; text-align: center; } .diagnosticEmpty strong { color: #cbd8ea; } +.taskflowFramePane, .taskflowRuntimePane { display: grid; align-content: start; gap: 14px; } +.taskflowCaptureBar, .taskflowSelectors { display: flex; align-items: center; flex-wrap: wrap; gap: 10px; padding: 11px 12px; border: 1px solid #29435e; border-radius: 9px; background: #0c1a2a; } +.taskflowCaptureBar label, .taskflowSelectors label { display: flex; align-items: center; gap: 7px; color: #91a5c0; font-size: 11px; } +.taskflowCaptureBar input, .taskflowSelectors select { min-width: 74px; padding: 7px 9px; color: #dce8f8; border: 1px solid #304664; border-radius: 7px; background: #101c2d; } +.taskflowCaptureBar button, .taskflowRuntimeHeader button { padding: 8px 12px; color: #062019; border: 1px solid #5ce4c2; border-radius: 7px; background: #5ce4c2; cursor: pointer; } +.taskflowCaptureBar button:disabled { opacity: .45; cursor: default; } +.taskflowCaptureBar span, .taskflowSelectors > span { color: #71839e; font-size: 11px; } +.taskflowError { margin: 0; padding: 10px 12px; color: #ff9bae; border: 1px solid #71334a; border-radius: 8px; background: #27101a; } +.taskflowGraphSummary, .taskflowRuntimeSummary { display: grid; grid-template-columns: repeat(auto-fit, minmax(130px, 1fr)); gap: 8px; margin: 0; } +.taskflowGraphSummary > div, .taskflowRuntimeSummary > div { min-width: 0; padding: 10px; border: 1px solid #213653; border-radius: 9px; background: #0a1422; } +.taskflowGraphSummary dt, .taskflowRuntimeSummary dt { color: #71839e; font-size: 10px; } +.taskflowGraphSummary dd, .taskflowRuntimeSummary dd { overflow: hidden; margin: 5px 0 0; color: #5ce4c2; font: 700 13px/1.25 ui-monospace, monospace; text-overflow: ellipsis; white-space: nowrap; } +.taskflowDag { width: 100%; height: 430px; min-height: 300px; overflow: hidden; border: 1px solid #213653; border-radius: 10px; background: #07101c; } +.taskflowNode { width: 230px !important; min-height: 92px; padding: 10px !important; color: #a7bad2 !important; border: 1px solid #304766 !important; border-radius: 9px !important; background: #0e1a2b !important; box-shadow: 0 8px 22px #0006; } +.taskflowNodeExecuted { border-color: #2c8e78 !important; background: #0c201e !important; } +.taskflowNodeWaiting { border-color: #b9833e !important; background: #251b10 !important; } +.taskflowNodeLabel { display: grid; gap: 4px; min-width: 0; text-align: left; } +.taskflowNodeLabel strong { overflow: hidden; color: #e1ecf9; font-size: 12px; text-overflow: ellipsis; white-space: nowrap; } +.taskflowNodeLabel code { overflow: hidden; color: #6f89aa; font: 9px/1.2 ui-monospace, monospace; text-overflow: ellipsis; white-space: nowrap; } +.taskflowNodeLabel span { color: #91a5c0; font: 9px/1.25 ui-monospace, monospace; } +.react-flow__edge-path { stroke: #557594; stroke-width: 1.5; } +.react-flow__controls button { color: #dce8f8; border-color: #29435e; background: #101d2f; } +.react-flow__minimap { border: 1px solid #29435e; background: #0a1422; } +.taskflowTimeline, .taskflowRuntimeSection { overflow: hidden; border: 1px solid #213653; border-radius: 10px; background: #0a1422; } +.taskflowTimeline > header, .taskflowRuntimeSection > header { display: flex; align-items: center; justify-content: space-between; gap: 12px; padding: 11px 13px; border-bottom: 1px solid #1d304a; background: #101d2f; } +.taskflowTimeline > header span, .taskflowRuntimeSection > header span { color: #71839e; font-size: 10px; } +.taskflowTimelineRows { display: grid; overflow-x: auto; padding: 9px; } +.taskflowTimelineRow { display: grid; grid-template-columns: 34px minmax(115px, .5fr) minmax(260px, 1.5fr) 78px 88px; align-items: center; gap: 8px; min-width: 720px; padding: 6px; border-bottom: 1px solid #16263b; color: #91a5c0; font-size: 10px; } +.taskflowTimelineRow:last-child { border-bottom: 0; } +.taskflowTimelineRow > code { color: #5ce4c2; } +.taskflowTimelineRow > span { overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } +.taskflowTimelineRow > b { color: #dce8f8; font: 10px/1 ui-monospace, monospace; } +.taskflowTimelineRow > em { color: #be935b; font: 10px/1 ui-monospace, monospace; font-style: normal; } +.taskflowTimelineTrack { position: relative; height: 9px; overflow: hidden; border-radius: 99px; background: #09111d; } +.taskflowTimelineTrack i { position: absolute; top: 0; bottom: 0; border-radius: 99px; } +.taskflowWait { background: #a16d36; } +.taskflowRun { background: #51dabc; box-shadow: 0 0 7px #51dabc88; } +.taskflowRuntimeHeader { display: flex; align-items: center; justify-content: space-between; gap: 14px; padding: 17px 16px 14px; border-bottom: 1px solid #20314b; background: #0b1524; } +.taskflowRuntimeHeader h2 { margin: 4px 0 0; font-size: 21px; } +.taskflowRuntimeHeader p { margin: 5px 0 0; color: #71839e; font-size: 11px; } +.taskflowWorkerGrid { display: grid; grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); gap: 9px; padding: 11px; } +.taskflowWorkerGrid article { padding: 10px; border: 1px solid #1f334e; border-radius: 8px; background: #091321; } +.taskflowWorkerGrid article > header { display: flex; justify-content: space-between; color: #dce8f8; font-size: 11px; } +.taskflowWorkerGrid article > header span { color: #5ce4c2; font: 11px/1 ui-monospace, monospace; } +.taskflowUtilization { height: 6px; overflow: hidden; margin: 9px 0; border-radius: 99px; background: #18273b; } +.taskflowUtilization i { display: block; height: 100%; border-radius: inherit; background: linear-gradient(90deg, #3aa58d, #5ce4c2); } +.taskflowWorkerGrid dl { display: grid; grid-template-columns: 1fr 1fr; gap: 5px; margin: 0; } +.taskflowWorkerGrid dl > div { display: flex; justify-content: space-between; gap: 5px; color: #71839e; font-size: 9px; } +.taskflowWorkerGrid dd { margin: 0; color: #aec0d8; font-family: ui-monospace, monospace; } +.taskflowActiveTask { display: block; overflow: hidden; margin-top: 8px; color: #e0a65d; font: 9px/1.3 ui-monospace, monospace; text-overflow: ellipsis; white-space: nowrap; } +.taskflowTable { width: 100%; border-collapse: collapse; font-size: 11px; } +.taskflowTable th, .taskflowTable td { padding: 9px 11px; border-bottom: 1px solid #182a42; text-align: right; } +.taskflowTable th:first-child, .taskflowTable td:first-child { text-align: left; } +.taskflowTable th { color: #71839e; font-weight: 500; } +.taskflowTable td { color: #c7d5e7; font-family: ui-monospace, monospace; } + .componentCard { margin-bottom: 13px; overflow: hidden; border: 1px solid #213653; border-radius: 11px; background: #0a1422; } .componentCard > summary { display: flex; align-items: center; justify-content: space-between; gap: 12px; padding: 11px 13px; color: #dce8f8; background: #101d2f; cursor: pointer; }