diff --git a/AGENTS.md b/AGENTS.md index 2187f94..443ee7c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,5 +1,12 @@ ## 设计约束 +`` +除非特别说明不允许写任何同步等待代码。 +除非特别说明不允许写任何同步等待代码。 +除非特别说明不允许写任何同步等待代码。 +除非特别说明不允许写任何同步等待代码。 +`` + 我使用 里的这个环境 D:\ae\proj\Aethera\CMakePresets.json "toolchain/vs2022.json" ====================[ 构建 | Aethera_Kernel_check | vs2022_debug ]================ "C:\Program Files\JetBrains\CLion 2026.1\bin\cmake\win\x64\bin\cmake.exe" --build D: diff --git a/kernel/src/kernel/Task_Graph.cpp b/kernel/src/kernel/Task_Graph.cpp index ab786b1..e54e742 100644 --- a/kernel/src/kernel/Task_Graph.cpp +++ b/kernel/src/kernel/Task_Graph.cpp @@ -1,5 +1,7 @@ #include "Task_Graph.hpp" #include "Task_Graph_Internal.hpp" +#include +#include #include #include #include @@ -12,6 +14,7 @@ struct Task_Graph::Private { std::string name{}; /* 调用方提供的业务节点名。 */ std::string node_id{}; /* 本图内去重后的业务节点 ID。 */ std::shared_ptr child{}; /* 模块引用的子图;普通节点为空。 */ + std::vector> attributes{}; }; explicit Private(std::string graph_name) : taskflow(std::move(graph_name)) {} @@ -43,6 +46,7 @@ struct Task_Graph::Private { node.parent_node_id = std::string(parent); node.name = source.name; node.type = std::string(tf::to_string(source.task.type())); + node.attributes = source.attributes; source.task.for_each_predecessor([&](tf::Task value) { node.predecessors.push_back( static_cast(value.hash_value())); @@ -98,6 +102,22 @@ void Task_Node::precede(const Task_Node& after) const { d->graph->nodes[after.d->index].task); } +Task_Node& Task_Node::describe(std::string key, std::string value) { + if (!d || !d->graph || d->generation != d->graph->generation || + d->index >= d->graph->nodes.size()) + throw std::invalid_argument("Task graph node handle is stale"); + if (key.empty()) + throw std::invalid_argument("Task graph node attribute key is empty"); + auto& attributes = d->graph->nodes[d->index].attributes; + const auto found = std::ranges::find( + attributes, key, &std::pair::first); + if (found == attributes.end()) + attributes.emplace_back(std::move(key), std::move(value)); + else + found->second = std::move(value); + return *this; +} + Task_Graph::Task_Graph(std::string graph_name) : d(std::make_shared(std::move(graph_name))) {} Task_Graph::~Task_Graph() = default; @@ -122,7 +142,20 @@ Task_Graph& Task_Graph::operator=(Task_Graph&& other) noexcept { Task_Node Task_Graph::add(std::string task_name, std::function work) { if (!work) throw std::invalid_argument("Task graph work is empty"); - auto task = d->taskflow.emplace(std::move(work)); + auto task = d->taskflow.emplace( + [name = task_name, work = std::move(work)] { + try { work(); } + catch (const std::exception& error) { + std::fprintf(stderr, "Taskflow node %s failed: %s\n", + name.c_str(), error.what()); + throw; + } + catch (...) { + std::fprintf(stderr, "Taskflow node %s failed: non-standard exception\n", + name.c_str()); + throw; + } + }); const auto node_id = d->unique_node_id(task_name); task.name(node_id); d->nodes.push_back({task, std::move(task_name), node_id, nullptr}); @@ -133,7 +166,20 @@ Task_Node Task_Graph::add(std::string task_name, Task_Node Task_Graph::add_condition(std::string task_name, std::function work) { if (!work) throw std::invalid_argument("Task graph condition is empty"); - auto task = d->taskflow.emplace(std::move(work)); + auto task = d->taskflow.emplace( + [name = task_name, work = std::move(work)] { + try { return work(); } + catch (const std::exception& error) { + std::fprintf(stderr, "Taskflow condition %s failed: %s\n", + name.c_str(), error.what()); + throw; + } + catch (...) { + std::fprintf(stderr, "Taskflow condition %s failed: non-standard exception\n", + name.c_str()); + throw; + } + }); const auto node_id = d->unique_node_id(task_name); task.name(node_id); d->nodes.push_back({task, std::move(task_name), node_id, nullptr}); diff --git a/kernel/src/kernel/Task_Graph.hpp b/kernel/src/kernel/Task_Graph.hpp index f4c2425..9823378 100644 --- a/kernel/src/kernel/Task_Graph.hpp +++ b/kernel/src/kernel/Task_Graph.hpp @@ -19,6 +19,7 @@ public: Task_Node& operator=(Task_Node&&) noexcept; /* 建立本节点到 after 的有向依赖;两个节点必须属于同一张图。 */ void precede(const Task_Node& after) const; + Task_Node& describe(std::string key, std::string value); private: friend class Task_Graph; struct Private; diff --git a/kernel/src/kernel/Taskflow_Frame_Access.hpp b/kernel/src/kernel/Taskflow_Frame_Access.hpp index 56a5c9c..0b680c2 100644 --- a/kernel/src/kernel/Taskflow_Frame_Access.hpp +++ b/kernel/src/kernel/Taskflow_Frame_Access.hpp @@ -18,10 +18,16 @@ struct Taskflow_Frame_Access { static void release_writer(Render_Frame& frame) noexcept; [[nodiscard]] static bool contains_task(const Render_Frame& frame, std::uint64_t native_id) noexcept; - static void append_task(Render_Frame& frame, std::size_t worker, - std::uint64_t native_id, std::size_t queue_size, - std::size_t queue_capacity, Clock::time_point started, - Clock::time_point finished); + [[nodiscard]] static std::size_t append_task( + Render_Frame& frame, std::size_t worker, std::uint64_t native_id, + std::size_t queue_size, std::size_t queue_capacity, + Clock::time_point entered, Clock::time_point started, + Clock::time_point finished, std::uint64_t cpu_entered_ns, + std::uint64_t cpu_started_ns, std::uint64_t cpu_finished_ns); + static void finish_task_observer( + Render_Frame& frame, std::size_t worker, std::size_t task, + Clock::time_point completed, std::uint64_t cpu_finished_ns, + std::uint64_t cpu_completed_ns) noexcept; [[nodiscard]] static Taskflow_Graph_Token begin_graph( Render_Frame& frame, Task_Graph& graph, std::string_view stage); static void finish_graph(Taskflow_Graph_Token token) noexcept; diff --git a/kernel/src/kernel/frame.cpp b/kernel/src/kernel/frame.cpp index 45f68ac..ce398e9 100644 --- a/kernel/src/kernel/frame.cpp +++ b/kernel/src/kernel/frame.cpp @@ -6,12 +6,14 @@ #include #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); @@ -174,7 +176,28 @@ Frame_Statistics_Sample Render_Frame::statistics(Frame_Dimension dimension) cons if (dimension == Frame_Dimension::two_dimensional) { result.set(Frame_Statistic::pipeline_2d_event_ms, take(event_time)); result.set(Frame_Statistic::pipeline_2d_prepare_ms, take(prepare_time)); - result.set(Frame_Statistic::pipeline_2d_paint_ms, take(paint_time)); + result.set(Frame_Statistic::pipeline_2d_paint_ms, paint_time); + double paint_remaining = paint_time; + const auto take_paint = [&](Frame_Trace_Marker first, + Frame_Trace_Marker last) { + const double value = std::min(paint_remaining, interval(first, last)); + paint_remaining -= value; + return take(value); + }; + result.set(Frame_Statistic::pipeline_2d_frame_target_ms, take_paint( + Frame_Trace_Marker::paint_frame_target_started, + Frame_Trace_Marker::paint_frame_target_finished)); + result.set(Frame_Statistic::pipeline_2d_background_ms, take_paint( + Frame_Trace_Marker::paint_background_started, + Frame_Trace_Marker::paint_background_finished)); + result.set(Frame_Statistic::pipeline_2d_cache_targets_ms, take_paint( + Frame_Trace_Marker::paint_cache_targets_started, + Frame_Trace_Marker::paint_cache_targets_finished)); + result.set(Frame_Statistic::pipeline_2d_taskflow_ms, take_paint( + Frame_Trace_Marker::paint_taskflow_started, + Frame_Trace_Marker::paint_taskflow_finished)); + result.set(Frame_Statistic::pipeline_2d_paint_coordination_ms, + take(paint_remaining)); result.set(Frame_Statistic::pipeline_2d_scene_coordination_ms, take(std::max(0.0, scene_time - event_time - prepare_time - paint_time))); result.set(Frame_Statistic::pipeline_2d_callback_ms, take(interval( @@ -255,6 +278,12 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { result.identity = d->identity; result.created_time_unix_ns = d->created_time_unix_ns; result.worker_count = d->taskflow_workers.size(); + for (std::size_t index = 0; index < marker_count; ++index) { + const auto encoded = d->markers[index].load(std::memory_order_acquire); + if (!encoded) continue; + result.markers.push_back(Frame_Trace_Point{ + static_cast(index), decode_present_value(encoded)}); + } { std::lock_guard lock(d->taskflow_graph_mutex); result.graphs = d->taskflow_graphs; @@ -271,24 +300,80 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { }); }); - std::unordered_map finished; + struct Node_Context { + const Taskflow_Graph_Trace* graph{}; + const Taskflow_Graph_Trace::Node* node{}; + }; + std::unordered_map nodes; + std::unordered_map native_id_by_node_id; + std::unordered_map> children_by_parent_id; + for (const auto& graph : result.graphs) { + for (const auto& node : graph.nodes) { + nodes.try_emplace(node.native_id, Node_Context{&graph, &node}); + native_id_by_node_id.try_emplace(node.node_id, node.native_id); + if (!node.parent_node_id.empty()) + children_by_parent_id[node.parent_node_id].push_back(node.native_id); + } + } + std::unordered_map observed_finished; + for (const auto& task : result.tasks) + observed_finished[task.native_id] = std::max( + observed_finished[task.native_id], task.completed_ms); + + std::unordered_map ready_cache; + std::unordered_map finished_cache; + std::unordered_set resolving_ready; + std::unordered_set resolving_finished; + std::function ready_of; + std::function finished_of; + ready_of = [&](std::uint64_t native_id) -> double { + if (const auto cached = ready_cache.find(native_id); + cached != ready_cache.end()) return cached->second; + const auto context = nodes.find(native_id); + if (context == nodes.end() || !resolving_ready.insert(native_id).second) + return 0.0; + double ready = context->second.graph->submitted_ms; + const auto& node = *context->second.node; + if (!node.parent_node_id.empty()) { + if (const auto parent = native_id_by_node_id.find(node.parent_node_id); + parent != native_id_by_node_id.end()) + ready = std::max(ready, ready_of(parent->second)); + } + for (const auto predecessor : node.predecessors) + ready = std::max(ready, finished_of(predecessor)); + resolving_ready.erase(native_id); + ready_cache.emplace(native_id, ready); + return ready; + }; + finished_of = [&](std::uint64_t native_id) -> double { + if (const auto cached = finished_cache.find(native_id); + cached != finished_cache.end()) return cached->second; + if (!resolving_finished.insert(native_id).second) return ready_of(native_id); + double finished = ready_of(native_id); + if (const auto observed = observed_finished.find(native_id); + observed != observed_finished.end()) + finished = std::max(finished, observed->second); + if (const auto context = nodes.find(native_id); context != nodes.end()) { + if (const auto children = children_by_parent_id.find( + context->second.node->node_id); + children != children_by_parent_id.end()) + for (const auto child : children->second) + finished = std::max(finished, finished_of(child)); + } + resolving_finished.erase(native_id); + finished_cache.emplace(native_id, finished); + return finished; + }; + + std::unordered_map previous_execution_finished; for (auto& task : result.tasks) { - double ready{}; - bool matched_graph{}; - for (const auto& graph : result.graphs) - for (const auto& node : graph.nodes) - if (node.native_id == task.native_id) { - if (!matched_graph || graph.submitted_ms < ready) - ready = graph.submitted_ms; - matched_graph = true; - for (const auto predecessor : node.predecessors) - if (const auto found = finished.find(predecessor); - found != finished.end()) - ready = std::max(ready, found->second); - } + double ready = ready_of(task.native_id); + if (const auto previous = previous_execution_finished.find(task.native_id); + previous != previous_execution_finished.end()) + ready = std::max(ready, previous->second); task.ready_ms = ready; - task.queue_wait_ms = std::max(0.0, task.started_ms - ready); - finished[task.native_id] = std::max(finished[task.native_id], task.finished_ms); + task.queue_wait_ms = std::max(0.0, task.entered_ms - ready); + previous_execution_finished[task.native_id] = task.completed_ms; } return result; } @@ -344,12 +429,15 @@ bool detail::Taskflow_Frame_Access::contains_task( } } -void detail::Taskflow_Frame_Access::append_task( +std::size_t detail::Taskflow_Frame_Access::append_task( Render_Frame& frame, std::size_t worker, std::uint64_t native_id, std::size_t queue_size, std::size_t queue_capacity, - Clock::time_point started, Clock::time_point finished) { + Clock::time_point entered, Clock::time_point started, + Clock::time_point finished, std::uint64_t cpu_entered_ns, + std::uint64_t cpu_started_ns, std::uint64_t cpu_finished_ns) { auto& data = *frame.d; - if (worker >= data.taskflow_workers.size()) return; + if (worker >= data.taskflow_workers.size()) + return std::numeric_limits::max(); const auto elapsed_ms = [&](Clock::time_point value) { return std::chrono::duration(value - data.created_at).count(); }; @@ -358,10 +446,37 @@ void detail::Taskflow_Frame_Access::append_task( trace.worker_id = worker; trace.worker_queue_size = queue_size; trace.worker_queue_capacity = queue_capacity; + trace.entered_ms = elapsed_ms(entered); trace.started_ms = elapsed_ms(started); trace.finished_ms = elapsed_ms(finished); + trace.completed_ms = trace.finished_ms; trace.duration_ms = std::max(0.0, trace.finished_ms - trace.started_ms); + trace.cpu_duration_ms = cpu_finished_ns >= cpu_started_ns + ? static_cast(cpu_finished_ns - cpu_started_ns) / 1'000'000.0 + : 0.0; + trace.observer_entry_ms = std::max(0.0, trace.started_ms - trace.entered_ms); + trace.observer_entry_cpu_ms = cpu_started_ns >= cpu_entered_ns + ? static_cast(cpu_started_ns - cpu_entered_ns) / 1'000'000.0 + : 0.0; data.taskflow_workers[worker].tasks.push_back(std::move(trace)); + return data.taskflow_workers[worker].tasks.size() - 1; +} + +void detail::Taskflow_Frame_Access::finish_task_observer( + Render_Frame& frame, std::size_t worker, std::size_t task, + Clock::time_point completed, std::uint64_t cpu_finished_ns, + std::uint64_t cpu_completed_ns) noexcept { + auto& data = *frame.d; + if (worker >= data.taskflow_workers.size() || + task >= data.taskflow_workers[worker].tasks.size()) return; + auto& trace = data.taskflow_workers[worker].tasks[task]; + trace.completed_ms = std::chrono::duration( + completed - data.created_at).count(); + trace.observer_exit_ms = std::max( + 0.0, trace.completed_ms - trace.finished_ms); + trace.observer_exit_cpu_ms = cpu_completed_ns >= cpu_finished_ns + ? static_cast(cpu_completed_ns - cpu_finished_ns) / 1'000'000.0 + : 0.0; } detail::Taskflow_Graph_Token detail::Taskflow_Frame_Access::begin_graph( diff --git a/kernel/src/kernel/frame.hpp b/kernel/src/kernel/frame.hpp index 52cca73..d420b26 100644 --- a/kernel/src/kernel/frame.hpp +++ b/kernel/src/kernel/frame.hpp @@ -17,6 +17,14 @@ enum class Frame_Trace_Marker : std::uint8_t { prepare_started, prepare_finished, paint_started, + paint_frame_target_started, + paint_frame_target_finished, + paint_background_started, + paint_background_finished, + paint_cache_targets_started, + paint_cache_targets_finished, + paint_taskflow_started, + paint_taskflow_finished, paint_finished, backend_queue_entered, backend_prepare_started, @@ -73,6 +81,7 @@ struct Taskflow_Graph_Trace { std::string type{}; /* Taskflow 原生 TaskType。 */ std::vector predecessors{}; /* 原生直接前驱 hash。 */ std::vector successors{}; /* 原生直接后继 hash。 */ + std::vector> attributes{}; }; std::vector nodes{}; /* DAG 构造时登记的节点与依赖元信息。 */ double submitted_ms{}; /* 相对帧创建时刻的 run 提交时间。 */ @@ -84,16 +93,24 @@ struct Taskflow_Task_Trace { std::size_t worker_id{}; /* 执行该任务的 Executor worker。 */ std::size_t worker_queue_size{}; /* on_entry 时的原生 worker queue_size。 */ std::size_t worker_queue_capacity{}; /* on_entry 时的原生 worker queue_capacity。 */ - double started_ms{}; /* 相对帧创建时刻的 on_entry 时间。 */ - double finished_ms{}; /* 相对帧创建时刻的 on_exit 时间。 */ - double duration_ms{}; /* 原生 on_entry 到 on_exit 的持续时间。 */ + double entered_ms{}; /* Observer on_entry 开始,即 Executor 真正选中任务的时间。 */ + double started_ms{}; /* Observer on_entry 返回前,即任务体即将执行的时间。 */ + double finished_ms{}; /* Observer on_exit 进入,即任务体已经结束的时间。 */ + double completed_ms{}; /* Observer on_exit 与按帧追踪写入全部结束的时间。 */ + double duration_ms{}; /* 仅任务体 started 到 finished 的持续时间。 */ + double cpu_duration_ms{}; /* 任务体在当前 worker 线程上实际消耗的 CPU 时间。 */ + double observer_entry_ms{}; /* on_entry 诊断本身的耗时。 */ + double observer_exit_ms{}; /* on_exit 诊断与按帧追踪写入的耗时。 */ + double observer_entry_cpu_ms{}; /* on_entry 诊断实际消耗的 worker CPU 时间。 */ + double observer_exit_cpu_ms{}; /* on_exit 诊断实际消耗的 worker CPU 时间。 */ double ready_ms{}; /* 前驱完成或根 run 提交后的估算就绪时间。 */ - double queue_wait_ms{}; /* ready 到 on_entry 的估算 Executor 排队时间。 */ + double queue_wait_ms{}; /* ready 到 entered 的估算 Executor 排队时间。 */ }; struct Taskflow_Frame_Trace { Frame_Identity identity{}; /* 该执行图所属逻辑渲染帧。 */ std::uint64_t created_time_unix_ns{}; /* 帧创建 Unix 时间,单位纳秒。 */ std::size_t worker_count{}; /* 捕获时全局 Executor 的 worker 数。 */ + std::vector markers{}; /* 与该帧 DAG 共用时间原点的原始流水线时间点。 */ std::vector graphs{}; /* 本帧主动执行的业务 DAG 元信息。 */ std::vector tasks{}; /* 本帧窗口内原生 Observer 完成的任务执行。 */ }; diff --git a/kernel/src/kernel/frame_statistics.hpp b/kernel/src/kernel/frame_statistics.hpp index 1af50db..9c77201 100644 --- a/kernel/src/kernel/frame_statistics.hpp +++ b/kernel/src/kernel/frame_statistics.hpp @@ -36,6 +36,11 @@ enum class Frame_Statistic : std::uint8_t { pipeline_2d_event_ms, pipeline_2d_prepare_ms, pipeline_2d_paint_ms, + pipeline_2d_frame_target_ms, + pipeline_2d_background_ms, + pipeline_2d_cache_targets_ms, + pipeline_2d_taskflow_ms, + pipeline_2d_paint_coordination_ms, pipeline_2d_scene_coordination_ms, pipeline_2d_callback_ms, pipeline_2d_frame_handoff_ms, diff --git a/kernel/src/kernel/render_common.cpp b/kernel/src/kernel/render_common.cpp index 098186b..a5a5569 100644 --- a/kernel/src/kernel/render_common.cpp +++ b/kernel/src/kernel/render_common.cpp @@ -5,10 +5,18 @@ #include #include #include +#include #include #include #include #include +#if defined(_WIN32) +#define WIN32_LEAN_AND_MEAN +#define NOMINMAX +#include +#elif defined(CLOCK_THREAD_CPUTIME_ID) +#include +#endif namespace aethera { namespace { tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept { @@ -16,6 +24,28 @@ tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept { detail::Task_Graph_Access::native_storage(graph)); } +std::uint64_t current_thread_cpu_ns() noexcept { +#if defined(_WIN32) + FILETIME created{}, exited{}, kernel{}, user{}; + if (!GetThreadTimes(GetCurrentThread(), &created, &exited, &kernel, &user)) + return 0; + const auto ticks = [](FILETIME value) noexcept { + ULARGE_INTEGER result{}; + result.LowPart = value.dwLowDateTime; + result.HighPart = value.dwHighDateTime; + return result.QuadPart; + }; + return (ticks(kernel) + ticks(user)) * 100ULL; +#elif defined(CLOCK_THREAD_CPUTIME_ID) + timespec value{}; + if (clock_gettime(CLOCK_THREAD_CPUTIME_ID, &value) != 0) return 0; + return static_cast(value.tv_sec) * 1'000'000'000ULL + + static_cast(value.tv_nsec); +#else + return 0; +#endif +} + class Task_Observer : public tf::ObserverInterface { private: using Clock = std::chrono::steady_clock; @@ -40,7 +70,10 @@ private: std::atomic_uint64_t max_task_time_ns{}; }; struct Start_Record { - Clock::time_point started{}; /* Observer on_entry 时间。 */ + Clock::time_point entered{}; /* Observer on_entry 进入时间。 */ + Clock::time_point started{}; /* on_entry 完成、任务体即将执行的时间。 */ + std::uint64_t cpu_entered_ns{}; /* on_entry 进入时的 worker CPU 时间。 */ + std::uint64_t cpu_started_ns{}; /* 任务体开始前的 worker CPU 时间。 */ Render_Frame* frame{}; /* 进入任务时唯一活动的按帧捕获。 */ std::size_t queue_size{}; /* 进入任务时 worker 队列深度。 */ std::size_t queue_capacity{}; /* 进入任务时 worker 队列容量。 */ @@ -110,7 +143,8 @@ public: update_max(peak_active_workers, active); } worker_starts.push_back(Start_Record{ - now, nullptr, worker.queue_size(), worker.queue_capacity(), + now, {}, current_thread_cpu_ns(), 0, nullptr, + worker.queue_size(), worker.queue_capacity(), static_cast(task.hash_value()), task.type()}); auto* frame = trace_frame.load(std::memory_order_acquire); /* @@ -148,13 +182,18 @@ public: update_max(max_weak_dependencies, task.num_weak_dependencies()); if (!task.name().empty()) named_tasks.fetch_add(1, std::memory_order_relaxed); update_first(first_task_time_ns, clock_ns(now)); + worker_starts.back().cpu_started_ns = current_thread_cpu_ns(); + worker_starts.back().started = Clock::now(); } void on_exit(tf::WorkerView worker, tf::TaskView task) override { - auto now = Clock::now(); + const auto finished = Clock::now(); + const auto cpu_finished_ns = current_thread_cpu_ns(); auto& worker_starts = starts[worker.id()]; auto start = worker_starts.back(); worker_starts.pop_back(); - auto elapsed = static_cast(std::chrono::duration_cast(now - start.started).count()); + const auto elapsed = static_cast( + std::chrono::duration_cast( + finished - start.started).count()); total_execution_time_ns.fetch_add(elapsed, std::memory_order_relaxed); completed_tasks.fetch_add(1, std::memory_order_relaxed); auto& worker_state = worker_statistics[worker.id()]; @@ -180,8 +219,33 @@ public: std::memory_order_release); } active_tasks.fetch_sub(1, std::memory_order_relaxed); + std::optional trace_task; + if (start.frame) { + try { + trace_task = detail::Taskflow_Frame_Access::append_task( + *start.frame, worker.id(), + static_cast(task.hash_value()), + start.queue_size, start.queue_capacity, start.entered, + start.started, finished, start.cpu_entered_ns, + start.cpu_started_ns, cpu_finished_ns); + } + catch (...) { + /* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */ + } + } + const auto cpu_completed_ns = current_thread_cpu_ns(); + const auto completed = Clock::now(); + if (start.frame) { + if (trace_task) + detail::Taskflow_Frame_Access::finish_task_observer( + *start.frame, worker.id(), *trace_task, completed, + cpu_finished_ns, cpu_completed_ns); + detail::Taskflow_Frame_Access::release_writer(*start.frame); + } if (worker_starts.empty()) { - auto busy = static_cast(std::chrono::duration_cast(now - worker_busy_starts[worker.id()]).count()); + auto busy = static_cast( + std::chrono::duration_cast( + completed - worker_busy_starts[worker.id()]).count()); worker_busy_time_ns.fetch_add(busy, std::memory_order_relaxed); worker_state.busy_time_ns.fetch_add(busy, std::memory_order_relaxed); active_workers.fetch_sub(1, std::memory_order_relaxed); @@ -195,27 +259,15 @@ public: else { const auto& parent = worker_starts.back(); worker_state.active_task_hash.store(parent.native_id, std::memory_order_relaxed); - worker_state.active_task_started_ns.store(clock_ns(parent.started), + worker_state.active_task_started_ns.store(clock_ns(parent.entered), std::memory_order_relaxed); worker_state.active_task_type.store(parent.type, std::memory_order_relaxed); worker_state.current_queue_size.store(parent.queue_size, std::memory_order_relaxed); worker_state.current_queue_capacity.store(parent.queue_capacity, - std::memory_order_relaxed); - } - update_max(last_task_time_ns, clock_ns(now)); - if (start.frame) { - try { - detail::Taskflow_Frame_Access::append_task( - *start.frame, worker.id(), - static_cast(task.hash_value()), - start.queue_size, start.queue_capacity, start.started, now); - } - catch (...) { - /* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */ - } - detail::Taskflow_Frame_Access::release_writer(*start.frame); + std::memory_order_relaxed); } + update_max(last_task_time_ns, clock_ns(completed)); } bool begin_trace(Render_Frame& frame, std::size_t workers) { std::unique_lock guard(trace_mutex); diff --git a/kernel/src/kernel/scene.ipp b/kernel/src/kernel/scene.ipp index e36bc90..42b1a40 100644 --- a/kernel/src/kernel/scene.ipp +++ b/kernel/src/kernel/scene.ipp @@ -164,8 +164,6 @@ void Scene::Private::after_advance(Object* object, if (dispatch->prepare.run) dispatch->prepare.run(root); }); } - auto prepare_extension = taskflow.compose(task_prefix + ".extension", - data->prepare_extension); auto prepare_done = taskflow.add(task_prefix + ".complete", [dispatch, root] { auto& state = *dispatch->state.pending(root); if (state.prepare_executed) { @@ -179,8 +177,15 @@ void Scene::Private::after_advance(Object* object, }); prepare_if.precede(prepare_run); prepare_if.precede(prepare_done); - prepare_run.precede(prepare_extension); - prepare_extension.precede(prepare_done); + if (data->prepare_extension.empty()) { + prepare_run.precede(prepare_done); + } + else { + auto prepare_extension = taskflow.compose( + task_prefix + ".extension", data->prepare_extension); + prepare_run.precede(prepare_extension); + prepare_extension.precede(prepare_done); + } stage_tasks.emplace(root, Stage_Tasks{prepare_if, prepare_done}); } auto connect_dependencies = [&](const auto& dependency_graph) { diff --git a/kernel/src/test/render_test.cpp b/kernel/src/test/render_test.cpp index ebdfae2..62e6809 100644 --- a/kernel/src/test/render_test.cpp +++ b/kernel/src/test/render_test.cpp @@ -222,21 +222,31 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) auto paint = child.add("paint.visual", [&] { completed.fetch_add(1); }); prepare.precede(paint); aethera::Task_Graph frame_graph{"test.frame"}; - frame_graph.compose("spectrum", child); + auto module_task = frame_graph.compose("spectrum", child); + auto publish = frame_graph.add("publish.state", [&] { completed.fetch_add(1); }); + module_task.precede(publish); aethera::Render_Frame frame{{41, 73}}; + frame.mark(aethera::Frame_Trace_Marker::paint_started); + frame.mark(aethera::Frame_Trace_Marker::paint_frame_target_started); + frame.mark(aethera::Frame_Trace_Marker::paint_frame_target_finished); + frame.mark(aethera::Frame_Trace_Marker::paint_finished); frame.request_taskflow_trace(); ASSERT_TRUE(aethera::detail::begin_taskflow_trace(frame)); aethera::detail::run_taskflow(frame_graph, frame, "test.scene.paint"); aethera::detail::finish_taskflow_trace(frame); - EXPECT_EQ(completed.load(), 2); + EXPECT_EQ(completed.load(), 3); const auto trace = frame.taskflow_trace(); ASSERT_EQ(trace.graphs.size(), 1u); EXPECT_EQ(trace.identity, (aethera::Frame_Identity{41, 73})); + EXPECT_TRUE(std::ranges::contains( + trace.markers, + aethera::Frame_Trace_Marker::paint_frame_target_started, + &aethera::Frame_Trace_Point::marker)); EXPECT_EQ(trace.graphs.front().stage, "test.scene.paint"); EXPECT_TRUE(trace.graphs.front().completed); - ASSERT_EQ(trace.graphs.front().nodes.size(), 3u); + ASSERT_EQ(trace.graphs.front().nodes.size(), 4u); const auto module = std::ranges::find( trace.graphs.front().nodes, "spectrum", &aethera::Taskflow_Graph_Trace::Node::name); @@ -254,4 +264,23 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) ASSERT_FALSE(trace.tasks.empty()); for (const auto& task : trace.tasks) EXPECT_TRUE(metadata_ids.contains(task.native_id)); + const auto child_paint = std::ranges::find( + trace.graphs.front().nodes, "paint.visual", + &aethera::Taskflow_Graph_Trace::Node::name); + const auto publish_node = std::ranges::find( + trace.graphs.front().nodes, "publish.state", + &aethera::Taskflow_Graph_Trace::Node::name); + ASSERT_NE(child_paint, trace.graphs.front().nodes.end()); + ASSERT_NE(publish_node, trace.graphs.front().nodes.end()); + const auto child_paint_execution = std::ranges::find( + trace.tasks, child_paint->native_id, + &aethera::Taskflow_Task_Trace::native_id); + const auto publish_execution = std::ranges::find( + trace.tasks, publish_node->native_id, + &aethera::Taskflow_Task_Trace::native_id); + ASSERT_NE(child_paint_execution, trace.tasks.end()); + ASSERT_NE(publish_execution, trace.tasks.end()); + EXPECT_NEAR(publish_execution->ready_ms, + child_paint_execution->finished_ms, 0.05); + EXPECT_LT(publish_execution->queue_wait_ms, 1.0); } diff --git a/render_2D/render_2D/scene/Render_Scene_2D.hpp b/render_2D/render_2D/scene/Render_Scene_2D.hpp index 8c92019..3381f8d 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.hpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.hpp @@ -40,7 +40,10 @@ struct Render_Scene_2D : Def; /* 向调用方拥有的帧合成一次;必须先安装回调。回调返回前的并发请求返回 frame_in_flight。 */ [[nodiscard]] std::expected render(Frame_2D* frame); - /* 安装完成帧回调;回调收到对应 render(frame) 的对象,返回时 Scene 才释放下一帧准入。 */ + /* + * 安装完成帧回调;回调收到对应 render(frame) 的对象,返回时 Scene 才释放下一帧准入。 + * Scene 和调用方拥有的 Frame 必须存活到回调返回;析构不会用同步等待隐藏错误的生命周期。 + */ void set_frame_callback(Frame_Callback callback); /* * 返回最终像素完成后、帧回调前执行的直接 Taskflow。 diff --git a/render_2D/render_2D/scene/Render_Scene_2D.ipp b/render_2D/render_2D/scene/Render_Scene_2D.ipp index bcb896c..c3e6740 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.ipp +++ b/render_2D/render_2D/scene/Render_Scene_2D.ipp @@ -31,6 +31,17 @@ std::expected, Dependency_Graph_Error> Render_Scene_2D:: auto attach_result = attach(scene.get()); if (!attach_result) return std::unexpected(attach_result.error()); } + /* + * Renderable attachment changes the Scene dependency graphs after the base + * Builder has performed its initial commit. Commit those attachments here, + * while no frame topology can be running, so the first render observes a + * fully built Prepare/Paint DAG instead of publishing a background-only + * bootstrap frame. + */ + scene->advance(); + /* after_advance writes derived DAG diagnostics into the State write side; + * publish that already-computed result before handing the Scene to Plot. */ + scene->advance(); return scene; } struct Render_Scene_2D::Private : Prev_Private { @@ -66,14 +77,18 @@ struct Render_Scene_2D::Private : Prev_Private { Frame_Statistics_Accumulator frame_statistics{}; /* Scene 内部增量计算;State 只发布定长统计结果。 */ Task_Graph completion_graph{"render_2d.completion"}; /* 最终像素完成后、发布回调前执行的外部续写图。 */ std::unique_ptr paint_taskflow{}; /* 仅由二维 Paint 图构建的执行图。 */ + std::unique_ptr frame_taskflow{}; + Frame_Callback active_callback{}; + bool active_trace{}; + std::vector active_prepare_executions{}; + std::chrono::steady_clock::time_point active_prepare_started{}; ~Private(); Frame_2D* active_frame{}; /* 当前同步 process 借用的外部帧;render 返回前清空。 */ Blend2D_Cache* frame_target{}; /* 当前 render(frame) 所属外部颜色层;调用返回后清空。 */ /* Impl CRTP 实现:在对象锁内执行 Kernel Scene,再按 Paint 图拓扑顺序合成颜色层。 */ std::vector paint_order{}; /* Paint 图当前拓扑序及 Scene 缓存分组结果。 */ std::vector cache_groups{}; /* 仅保存显式缓存根对应的执行分组。 */ - template - void process(Object* object, Callback&& callback) requires std::invocable; + template void ensure_frame_taskflow(Object* object); /* Def CRTP hook:二维 Paint 图结构变化后重建其独立 Taskflow。 */ template void after_advance(Object* object, Prop* pending_prop, State_Access pending_states, const Prop* current_prop, State_Access current_states); @@ -176,7 +191,6 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access 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) { @@ -190,8 +204,15 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access }); paint_if.precede(paint_run); paint_if.precede(paint_done); - paint_run.precede(paint_extension); - paint_extension.precede(paint_done); + if (data.paint_extension.empty()) { + paint_run.precede(paint_done); + } + else { + auto paint_extension = taskflow.compose( + task_prefix + ".extension", data.paint_extension); + 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) { @@ -254,63 +275,102 @@ void Render_Scene_2D::Private::prepare_paint_targets(Object*, Blend2D_Cache& fra node.object->template mark_dirty(); } } -template -void Render_Scene_2D::Private::process(Object* object, Callback&& callback) - requires std::invocable { - Frame_2D* frame_object = active_frame; - if (!frame_object) throw std::logic_error("2D Scene process has no external frame"); - const auto& state = static_cast(*static_cast(*this).current); - if (!state.view_active || state.viewport.empty()) return; - frame_object->mark(Frame_Trace_Marker::scene_render_started); - frame_object->mark(Frame_Trace_Marker::event_dispatch_started); - dispatch_events(object, state.viewport, frame_object->identity().sequence); - frame_object->mark(Frame_Trace_Marker::event_dispatch_finished); - object->template current_dependency_graph().for_each([](const Dependency_Graph::Node& node) { - if (!node.object->template take_dirty()) return; - node.object->template mark_dirty(); - node.object->template mark_dirty(); +template +void Render_Scene_2D::Private::ensure_frame_taskflow(Object* object) { + if (frame_taskflow) return; + if (!runtime->taskflow) + runtime->taskflow = std::make_unique("scene.prepare"); + if (!paint_taskflow) + paint_taskflow = std::make_unique("render_2d.paint"); + frame_taskflow = std::make_unique("render_2d.frame"); + auto& graph = *frame_taskflow; + auto begin = graph.add("scene.begin", [this, object] { + auto* frame = active_frame; + if (!frame) throw std::logic_error("2D frame DAG lost its active frame"); + const auto& prop = object->template read_prop(); + frame->mark(Frame_Trace_Marker::scene_render_started); + frame->mark(Frame_Trace_Marker::event_dispatch_started); + dispatch_events(object, prop.viewport, frame->identity().sequence); + frame->mark(Frame_Trace_Marker::event_dispatch_finished); + object->template current_dependency_graph().for_each( + [](const Dependency_Graph::Node& node) { + if (!node.object->template take_dirty()) return; + node.object->template mark_dirty(); + node.object->template mark_dirty(); + }); + active_prepare_executions.clear(); + object->template current_dependency_graph().for_each_bound( + [this](Renderable* renderable, Renderable::Private& data) { + bool execute = renderable->template dirty(); + if (data.dispatch->prepare.builder) + execute = execute || !data.prepare_graph_built || + data.dispatch->prepare.rebuild_predicate(renderable); + if (execute) active_prepare_executions.push_back(renderable); + }); + auto& private_data = static_cast(*this); + auto& state = static_cast(*private_data.state.pending); + state.taskflow_execution_time_ns = 0; + active_prepare_started = std::chrono::steady_clock::now(); + frame->mark(Frame_Trace_Marker::prepare_started); }); - std::vector prepare_executions; - object->template current_dependency_graph().for_each_bound([&](Renderable* renderable, Renderable::Private& data) { - bool execute = renderable->template dirty(); - if (data.dispatch->prepare.builder) - execute = execute || !data.prepare_graph_built || data.dispatch->prepare.rebuild_predicate(renderable); - if (execute) prepare_executions.push_back(renderable); - }); - Prev_Private::process(object, frame_object, [&](const Scene::Private::Result&) { + begin.describe("dimension", "2D").describe("stage", "event dispatch"); + auto prepare = graph.compose("scene.prepare", *runtime->taskflow); + prepare.describe("dimension", "2D").describe("stage", "data prepare"); + auto setup_paint = graph.add("scene.paint.setup", [this, object] { + auto* frame_object = active_frame; + if (!frame_object) throw std::logic_error("2D Paint lost its active frame"); + auto& private_data = static_cast(*this); + auto& scene_state = static_cast(*private_data.state.pending); + scene_state.taskflow_execution_time_ns = static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now() - active_prepare_started).count()); + scene_state.event_statistics = event_statistics.state(); + frame_object->mark(Frame_Trace_Marker::prepare_finished); frame_object->mark(Frame_Trace_Marker::paint_started); - for (auto* renderable : prepare_executions) renderable->template mark_dirty(); + for (auto* renderable : active_prepare_executions) + renderable->template mark_dirty(); auto& frame = detail::Frame_2D_Access::render_target(frame_object); - struct Frame_Target_Scope { - Blend2D_Cache*& target; /* Scene 当前帧目标槽位。 */ - Blend2D_Cache* previous{}; /* 嵌套调用前的目标;析构时恢复。 */ - ~Frame_Target_Scope() { target = previous; } - } frame_target_scope{frame_target, frame_target}; frame_target = &frame; - frame.ensure_size(state.viewport); + const auto& prop = object->template read_prop(); + frame_object->mark(Frame_Trace_Marker::paint_frame_target_started); + frame.ensure_size(prop.viewport); frame.clear(); + frame_object->mark(Frame_Trace_Marker::paint_frame_target_finished); + frame_object->mark(Frame_Trace_Marker::paint_background_started); { - detail::Painter painter(frame, state.viewport); - painter.rect({0.0, 0.0, static_cast(state.viewport.width), - static_cast(state.viewport.height)}, - Pen{.style = Line_Style::none}, Brush{state.background, Brush_Style::solid}); + detail::Painter painter(frame, prop.viewport); + painter.rect({0.0, 0.0, static_cast(prop.viewport.width), + static_cast(prop.viewport.height)}, + Pen{.style = Line_Style::none}, + Brush{prop.background, Brush_Style::solid}); } - prepare_paint_targets(object, frame, state.viewport); - 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); + frame_object->mark(Frame_Trace_Marker::paint_background_finished); + frame_object->mark(Frame_Trace_Marker::paint_cache_targets_started); + prepare_paint_targets(object, frame, prop.viewport); + frame_object->mark(Frame_Trace_Marker::paint_cache_targets_finished); + frame_object->mark(Frame_Trace_Marker::paint_taskflow_started); }); - std::invoke(std::forward(callback)); + setup_paint.describe("dimension", "2D").describe("stage", "frame target"); + auto paint = graph.compose("scene.paint", *paint_taskflow); + paint.describe("dimension", "2D").describe("backend", "Blend2D"); + auto finish_paint = graph.add("scene.paint.complete", [this] { + active_frame->mark(Frame_Trace_Marker::paint_taskflow_finished); + active_frame->mark(Frame_Trace_Marker::paint_finished); + active_frame->mark(Frame_Trace_Marker::scene_render_finished); + }); + finish_paint.describe("dimension", "2D").describe("stage", "pixel complete"); + auto completion = graph.compose("scene.completion", completion_graph); + completion.describe("owner", "scene").describe("stage", "frame consumers"); + begin.precede(prepare); + prepare.precede(setup_paint); + setup_paint.precede(paint); + paint.precede(finish_paint); + finish_paint.precede(completion); } + template std::expected Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { - static_cast(detail::Frame_2D_Access::render_target(frame)); Frame_Callback callback; { std::lock_guard lock(render_mutex); @@ -321,59 +381,88 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { frame_in_flight = true; callback = frame_callback; } - struct Admission_Scope { - Private& data; - ~Admission_Scope() { - std::lock_guard lock(data.render_mutex); - data.frame_in_flight = false; + const auto release_admission = [this] { + std::lock_guard lock(render_mutex); + frame_in_flight = false; + }; + try { + /* + * Plot writes properties and tagged input buffers before render(). + * Commit the complete Scene object here, before submitting the frame + * topology: this advances every attached Renderable, publishes the + * current input buffers, and rebuilds dirty Prepare/Paint graphs while + * no Taskflow is running. The completion callback only publishes + * results; it must not modify a graph that is still executing. + */ + object->advance(); + const auto& prop = + object->template read_prop(); + if (!prop.view_active) { + release_admission(); + return std::unexpected(Render_Result::view_inactive); } - } admission{*this}; + if (prop.viewport.empty()) { + release_admission(); + return std::unexpected(Render_Result::empty_viewport); + } + ensure_frame_taskflow(object); + active_frame = frame; + active_callback = callback; + } + catch (...) { + release_admission(); + throw; + } 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{}; /* 嵌套调用前的帧;析构时恢复。 */ - ~Active_Frame_Scope() { target = previous; } - } active_frame_scope{active_frame, active_frame}; - active_frame = frame; - bool completed{}; - object->process([&] { - completed = true; - frame->mark(Frame_Trace_Marker::scene_render_finished); - /* 帧准入直到 completion 图和最终回调都结束才释放,因此图运行期间外部帧不会被下一帧复用。 */ - 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) { + active_trace = aethera::detail::begin_taskflow_trace(*frame); + try { + auto completion = [this, object, frame] { + struct Release { + Private* data; + ~Release() { + data->active_frame = nullptr; + data->frame_target = nullptr; + data->active_callback = {}; + data->active_prepare_executions.clear(); + std::lock_guard lock(data->render_mutex); + data->frame_in_flight = false; + } + } release{this}; + if (active_trace) { + aethera::detail::finish_taskflow_trace(*frame); + active_trace = false; + } + frame->mark(Frame_Trace_Marker::callback_started); + active_callback(frame); + frame->mark(Frame_Trace_Marker::callback_finished); + frame->mark(Frame_Trace_Marker::frame_ready); + auto& private_data = static_cast(*this); + auto& scene_state = static_cast(*private_data.state.pending); + scene_state.frame_statistics = frame_statistics.submit( + *frame, Frame_Dimension::two_dimensional); + private_data.state.advance(); + object->template notify_state(); + }; + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow( + *frame_taskflow, *frame, "render_2d.frame", std::move(completion)); + else + aethera::detail::run_taskflow(*frame_taskflow, std::move(completion)); + } + catch (...) { + if (active_trace) { aethera::detail::finish_taskflow_trace(*frame); - trace_scope.active = false; + active_trace = false; } - frame->mark(Frame_Trace_Marker::callback_started); - callback(frame); - frame->mark(Frame_Trace_Marker::callback_finished); - frame->mark(Frame_Trace_Marker::frame_ready); - auto& private_data = static_cast(*this); - auto& scene_state = static_cast(*private_data.state.pending); - scene_state.frame_statistics = frame_statistics.submit( - *frame, Frame_Dimension::two_dimensional); - private_data.state.advance(); - object->template notify_state(); - }); - if (completed) return {}; - const auto& prop = object->template read_prop(); - return std::unexpected(prop.view_active ? Render_Result::empty_viewport - : Render_Result::view_inactive); + active_frame = nullptr; + frame_target = nullptr; + active_callback = {}; + active_prepare_executions.clear(); + std::lock_guard lock(render_mutex); + frame_in_flight = false; + throw; + } + return {}; } template void Render_Scene_2D::Private::dispatch_events(Object* object, Size viewport, diff --git a/render_3D/render_3D/detail/Async_Render_Backend.cpp b/render_3D/render_3D/detail/Async_Render_Backend.cpp index 372ed64..ef1096c 100644 --- a/render_3D/render_3D/detail/Async_Render_Backend.cpp +++ b/render_3D/render_3D/detail/Async_Render_Backend.cpp @@ -7,8 +7,6 @@ #include #include #include -#include -#include #include #include #include @@ -16,7 +14,6 @@ #include #include #include -#include #include namespace aethera::render_3d::detail { namespace { @@ -50,12 +47,30 @@ void record_datoviz_trace(Frame_3D* frame, const Datoviz_Frame_Trace& trace) { frame->record(Frame_Trace_Measurement::gpu_copy_ns, trace.gpu->copy_ns); frame->record(Frame_Trace_Measurement::gpu_total_ns, trace.gpu->total_ns); } +void append_exception_description(const std::exception& error, + std::string& result) { + if (!result.empty()) result += ": "; + result += error.what(); + const auto* nested = dynamic_cast(&error); + if (!nested) return; + try { + nested->rethrow_nested(); + } + catch (const std::exception& cause) { + append_exception_description(cause, result); + } + catch (...) { + result += ": non-standard nested exception"; + } +} std::string failure_description(const std::exception_ptr& failure) { try { if (failure) std::rethrow_exception(failure); } catch (const std::exception& error) { - return error.what(); + std::string result; + append_exception_description(error, result); + return result; } catch (...) { return "non-standard exception"; @@ -79,6 +94,7 @@ struct Async_Render_Backend::Implementation Datoviz_Visual_Backend::Pending_Frame backend_frame{}; /* 三缓冲目标对应的已录制提交。 */ std::vector completions{}; /* GPU 完成时一起结清的逻辑帧。 */ Extent extent{}; /* 失败时生成空结果所需的尺寸。 */ + std::optional completion_reservation{}; std::atomic_bool submitted{}; /* 区分提交前取消与 GPU 已拥有资源后的终止。 */ }; struct Gpu_Completion { @@ -86,15 +102,8 @@ struct Async_Render_Backend::Implementation std::optional result{}; /* fence 正常交付时的结果。 */ std::exception_ptr failure{}; /* 提交域或完成服务的 Unknown Failure。 */ }; - struct Stop {}; - using Command = std::variant; - std::shared_ptr render_domain; /* 同 GPU 唯一的轻量 Queue Submit 域。 */ std::unique_ptr backend{}; /* Builder build() 创建;随后由本 Scene 串行使用。 */ - std::optional deferred_submission{}; /* 三槽占满时仅保留的最新画面和历史回调。 */ - std::mutex command_mutex{}; /* 只保护本 Scene 的入队和 drain 所有权。 */ - std::deque commands{}; /* 不同 Scene 可并行,同一 Scene 保持严格串行。 */ - bool drain_scheduled{}; /* 是否已有 Taskflow worker 负责本 Scene。 */ std::mutex callback_mutex{}; /* 保护完成回调替换与 render 捕获。 */ Frame_Callback frame_callback{}; /* 后续 render 捕获的完成出口。 */ std::mutex failure_mutex{}; /* 保护跨线程传播的最后 Unknown Failure。 */ @@ -114,16 +123,11 @@ struct Async_Render_Backend::Implementation void stop() noexcept; void fail(std::exception_ptr value) noexcept; - void enqueue(Command command); - void drain() noexcept; - void execute(Command command) noexcept; - void finish_stop(); [[nodiscard]] Async_Render_Backend::Submit_Result render( Shared_Prepared_Visual_Batch visuals, Scene_3D_Parameters parameters, Frame_3D* frame, Scene::Event_Batch events); - void accept(Submission submission); - void merge_deferred(Submission submission); - void prepare(Submission submission); + void accept(Submission submission, std::shared_ptr pending); + void prepare(Submission submission, std::shared_ptr pending); void submit(std::shared_ptr pending); void finish(Gpu_Completion completion); void resolve(std::shared_ptr pending, @@ -153,7 +157,7 @@ void Async_Render_Backend::Implementation::stop() noexcept { available.store(false, std::memory_order_release); } try { - enqueue(Stop{}); + if (backend && backend->idle()) backend.reset(); } catch (...) { fail(std::current_exception()); } } @@ -173,36 +177,6 @@ void Async_Render_Backend::Implementation::fail( } catch (...) {} } -void Async_Render_Backend::Implementation::enqueue(Command command) { - bool schedule{}; - { - std::lock_guard lock(command_mutex); - commands.push_back(std::move(command)); - if (!drain_scheduled) { - drain_scheduled = true; - schedule = true; - } - } - if (!schedule) return; - auto self = shared_from_this(); - aethera::schedule_task("render_3d.backend.drain", - [self] { self->drain(); }); -} -void Async_Render_Backend::Implementation::drain() noexcept { - for (;;) { - std::optional command; - { - std::lock_guard lock(command_mutex); - if (commands.empty()) { - drain_scheduled = false; - break; - } - command.emplace(std::move(commands.front())); - commands.pop_front(); - } - execute(std::move(*command)); - } -} Async_Render_Backend::Submit_Result Async_Render_Backend::Implementation::render( Shared_Prepared_Visual_Batch visuals, Scene_3D_Parameters parameters, @@ -221,10 +195,40 @@ Async_Render_Backend::Implementation::render( callback = frame_callback; } frame->mark(Frame_Trace_Marker::backend_queue_entered); + auto pending = std::make_shared(); + pending->extent = parameters.viewport; + pending->completions = {Completion{frame, std::move(callback)}}; + auto self = shared_from_this(); + auto reservation = Gpu_Completion_Service::instance().prepare( + [self, pending](Gpu_Completion_Service::Result result) { + aethera::schedule_task("render_3d.backend.complete", + [self, pending, result] { + try { self->finish(Gpu_Completion{pending, result, {}}); } + catch (...) { self->fail(std::current_exception()); } + }); + }, + [self, pending](std::exception_ptr value) { + aethera::schedule_task("render_3d.backend.complete.failure", + [self, pending, value = std::move(value)]() mutable { + try { self->finish(Gpu_Completion{ + pending, std::nullopt, std::move(value)}); } + catch (...) { self->fail(std::current_exception()); } + }); + }, frame->identity().sequence == 1 || + frame->identity().sequence % backend_observation_period == 0); + if (!reservation) { + if (reservation.result == + Gpu_Completion_Service::Admission_Result::capacity_exhausted) + return Async_Render_Backend::Submit_Result::completion_capacity_exhausted; + return Async_Render_Backend::Submit_Result::backend_unavailable; + } + pending->completion_reservation.emplace(std::move(reservation.reservation)); try { - enqueue(Submission{std::move(visuals), std::move(parameters), - {Completion{frame, std::move(callback)}}, - std::move(events)}); + /* CPU 事件处理、Datoviz apply/emit/execute 与帧目标准备属于调用 Scene + * 的 frame Taskflow。这里先完成可并行的重工作,Render Domain 只接收 + * 已录制 Pending_Frame 的轻量 Vulkan submit。 */ + accept(Submission{std::move(visuals), std::move(parameters), {}, + std::move(events)}, std::move(pending)); } catch (...) { fail(std::current_exception()); @@ -232,50 +236,37 @@ Async_Render_Backend::Implementation::render( } return Async_Render_Backend::Submit_Result::queued; } -void Async_Render_Backend::Implementation::merge_deferred( - Submission submission) { - if (deferred_submission) { - submission.events.insert( - submission.events.begin(), - std::make_move_iterator(deferred_submission->events.begin()), - std::make_move_iterator(deferred_submission->events.end())); - submission.completions.insert( - submission.completions.begin(), - std::make_move_iterator(deferred_submission->completions.begin()), - std::make_move_iterator(deferred_submission->completions.end())); - } - deferred_submission = std::move(submission); -} -void Async_Render_Backend::Implementation::accept(Submission submission) { +void Async_Render_Backend::Implementation::accept( + Submission submission, std::shared_ptr pending) { if (stop_requested.load(std::memory_order_acquire) || !available.load(std::memory_order_acquire)) { - auto pending = std::make_shared(); - pending->extent = submission.parameters.viewport; - pending->completions = std::move(submission.completions); resolve(std::move(pending), std::nullopt); return; } if (backend && !backend->can_prepare(submission.parameters, *submission.visuals)) { - merge_deferred(std::move(submission)); + fail(std::make_exception_ptr(std::logic_error( + "3D Scene admitted a new frame before its previous frame completed"))); + resolve(std::move(pending), std::nullopt); return; } - prepare(std::move(submission)); + prepare(std::move(submission), std::move(pending)); } -void Async_Render_Backend::Implementation::prepare(Submission submission) { +void Async_Render_Backend::Implementation::prepare( + Submission submission, std::shared_ptr pending) { std::optional prepared; try { const auto sequence = - submission.completions.back().output->identity().sequence; + pending->completions.back().output->identity().sequence; for (const auto& event : submission.events) { if (event) event->mark_dispatch_started(sequence); dispatch(event, submission.parameters.viewport); } submission.events.clear(); - for (const auto& completion : submission.completions) + for (const auto& completion : pending->completions) completion.output->mark(Frame_Trace_Marker::backend_prepare_started); const bool readback = std::ranges::any_of( - submission.completions, [](const Completion& completion) { + pending->completions, [](const Completion& completion) { return completion.output->output() == Frame_3D_Output::pixels; }); const bool observe = sequence == 1 || @@ -285,25 +276,21 @@ void Async_Render_Backend::Implementation::prepare(Submission submission) { } catch (...) { fail(std::current_exception()); - auto pending = std::make_shared(); - pending->extent = submission.parameters.viewport; - pending->completions = std::move(submission.completions); resolve(std::move(pending), std::nullopt); return; } if (!prepared) { - merge_deferred(std::move(submission)); + fail(std::make_exception_ptr(std::logic_error( + "Datoviz did not provide a frame target after Scene admission"))); + resolve(std::move(pending), std::nullopt); return; } - for (const auto& completion : submission.completions) { + for (const auto& completion : pending->completions) { completion.output->mark(Frame_Trace_Marker::backend_prepare_finished); record_datoviz_trace(completion.output, prepared->trace); completion.output->mark(Frame_Trace_Marker::backend_submit_queued); } - auto pending = std::make_shared(); pending->backend_frame = std::move(*prepared); - pending->extent = submission.parameters.viewport; - pending->completions = std::move(submission.completions); try { submit(pending); } @@ -316,56 +303,38 @@ void Async_Render_Backend::Implementation::prepare(Submission submission) { void Async_Render_Backend::Implementation::submit( std::shared_ptr pending) { auto self = shared_from_this(); + /* Completion 槽位分配、回调所有权建立和服务登记没有 GPU 线程亲和性。 + * 它们在当前 Scene Taskflow worker 完成;唯一 Render Domain 只保留 + * vkQueueSubmit 与 fence 交付。Reservation 随 Pending 存活,提交前失败 + * 会由析构自动取消,不产生第二套生命周期状态。 */ const auto queued = render_domain->post( [self, pending] { for (const auto& completion : pending->completions) completion.output->mark(Frame_Trace_Marker::backend_queue_left); - auto reservation = Gpu_Completion_Service::instance().prepare( - [self, pending](Gpu_Completion_Service::Result result) { - try { - self->enqueue(Gpu_Completion{pending, result, {}}); - } - catch (...) { - self->fail(std::current_exception()); - } - }, - [self, pending](std::exception_ptr value) { - try { - self->enqueue(Gpu_Completion{pending, std::nullopt, - std::move(value)}); - } - catch (...) { - self->fail(std::current_exception()); - } - }, pending->backend_frame.trace.observed); - if (!reservation) { - self->enqueue(Gpu_Completion{ - pending, std::nullopt, - std::make_exception_ptr(std::runtime_error( - "GPU completion service stopped before submission"))}); - return; - } self->backend->submit(pending->backend_frame); pending->submitted.store(true, std::memory_order_release); for (const auto& completion : pending->completions) completion.output->mark(Frame_Trace_Marker::gpu_submitted); - reservation.reservation.watch(pending->backend_frame.device, - pending->backend_frame.fence); + pending->completion_reservation->watch( + pending->backend_frame.device, pending->backend_frame.fence); + pending->completion_reservation.reset(); }, [self, pending](std::exception_ptr value) { - try { - self->enqueue(Gpu_Completion{pending, std::nullopt, - std::move(value)}); - } - catch (...) { - self->fail(std::current_exception()); - } + aethera::schedule_task("render_3d.backend.complete.failure", + [self, pending, value = std::move(value)]() mutable { + try { self->finish(Gpu_Completion{ + pending, std::nullopt, std::move(value)}); } + catch (...) { self->fail(std::current_exception()); } + }); }); if (queued != Render_Domain::Post_Result::queued) { - enqueue(Gpu_Completion{ - std::move(pending), std::nullopt, - std::make_exception_ptr(std::runtime_error( - "render domain stopped before 3D queue submission"))}); + aethera::schedule_task("render_3d.backend.complete.failure", + [self, pending = std::move(pending)] { + self->finish(Gpu_Completion{ + pending, std::nullopt, + std::make_exception_ptr(std::runtime_error( + "render domain stopped before 3D queue submission"))}); + }); } } void Async_Render_Backend::Implementation::resolve( @@ -447,21 +416,8 @@ void Async_Render_Backend::Implementation::finish(Gpu_Completion completion) { std::move(completion.pending->backend_frame)); } resolve(std::move(completion.pending), std::move(completed)); - if (deferred_submission && - !stop_requested.load(std::memory_order_acquire) && - available.load(std::memory_order_acquire)) { - auto next = std::move(*deferred_submission); - deferred_submission.reset(); - prepare(std::move(next)); - } - else if (deferred_submission && - !available.load(std::memory_order_acquire)) { - auto pending = std::make_shared(); - pending->extent = deferred_submission->parameters.viewport; - pending->completions = std::move(deferred_submission->completions); - deferred_submission.reset(); - resolve(std::move(pending), std::nullopt); - } + if (stop_requested.load(std::memory_order_acquire) && backend && backend->idle()) + backend.reset(); } void Async_Render_Backend::Implementation::dispatch( const std::shared_ptr& event, Extent viewport) { @@ -504,44 +460,6 @@ void Async_Render_Backend::Implementation::dispatch( if ((wheel && pointer_valid) || pointer_valid || key) event->accept(); event->mark_dispatch_completed(); } -void Async_Render_Backend::Implementation::execute(Command command) noexcept { - try { - if (auto* submission = std::get_if(&command)) - accept(std::move(*submission)); - else if (auto* completion = std::get_if(&command)) - finish(std::move(*completion)); - else - finish_stop(); - } - catch (...) { - fail(std::current_exception()); - if (auto* submission = std::get_if(&command)) { - auto pending = std::make_shared(); - pending->extent = submission->parameters.viewport; - pending->completions = std::move(submission->completions); - try { resolve(std::move(pending), std::nullopt); } - catch (...) { fail(std::current_exception()); } - } - else if (auto* completion = std::get_if(&command)) { - try { resolve(std::move(completion->pending), std::nullopt); } - catch (...) { fail(std::current_exception()); } - } - } - if (stop_requested.load(std::memory_order_acquire)) finish_stop(); -} - -void Async_Render_Backend::Implementation::finish_stop() { - if (deferred_submission) { - auto pending = std::make_shared(); - pending->extent = deferred_submission->parameters.viewport; - pending->completions = std::move(deferred_submission->completions); - deferred_submission.reset(); - resolve(std::move(pending), std::nullopt); - } - if (backend && !backend->idle()) return; - backend.reset(); -} - Async_Render_Backend::Submit_Result Async_Render_Backend::render( Shared_Prepared_Visual_Batch visuals, Scene_3D_Parameters parameters, Frame_3D* frame, Scene::Event_Batch events) { diff --git a/render_3D/render_3D/detail/Async_Render_Backend.hpp b/render_3D/render_3D/detail/Async_Render_Backend.hpp index e9a2ee9..0efec4d 100644 --- a/render_3D/render_3D/detail/Async_Render_Backend.hpp +++ b/render_3D/render_3D/detail/Async_Render_Backend.hpp @@ -10,6 +10,7 @@ public: using Frame_Callback = std::function; enum class Submit_Result : std::uint8_t { queued, + completion_capacity_exhausted, backend_unavailable }; Async_Render_Backend(std::uint32_t gpu_index, bool validation_enabled, diff --git a/render_3D/render_3D/detail/Backend_Types.hpp b/render_3D/render_3D/detail/Backend_Types.hpp index 2e525cf..166e65c 100644 --- a/render_3D/render_3D/detail/Backend_Types.hpp +++ b/render_3D/render_3D/detail/Backend_Types.hpp @@ -5,6 +5,7 @@ #include "../camera/Camera.hpp" #include #include +#include #include #include namespace aethera::render_3d::detail { @@ -50,6 +51,18 @@ struct Prepared_Visual_Instance { using Prepared_Visual_Batch = std::vector; using Shared_Prepared_Visual_Batch = std::shared_ptr; +[[nodiscard]] inline bool prepared_visual_has_payload( + const Prepared_Visual& visual) noexcept { + return std::visit([](const auto& data) { return static_cast(data); }, + visual.data); +} +[[nodiscard]] inline bool prepared_visual_batch_complete( + const Prepared_Visual_Batch& visuals) noexcept { + return !visuals.empty() && + std::ranges::all_of(visuals, [](const Prepared_Visual_Instance& value) { + return prepared_visual_has_payload(value.visual); + }); +} struct Scene_3D_Parameters { Extent viewport{560, 320}; /* 离屏渲染目标尺寸,单位为像素。 */ Linear_Color clear_color{}; /* 每帧开始时写入的线性背景色。 */ diff --git a/render_3D/render_3D/detail/Datoviz_Visual_Backend.cpp b/render_3D/render_3D/detail/Datoviz_Visual_Backend.cpp index 3cbb583..98fc6d6 100644 --- a/render_3D/render_3D/detail/Datoviz_Visual_Backend.cpp +++ b/render_3D/render_3D/detail/Datoviz_Visual_Backend.cpp @@ -180,7 +180,26 @@ bool supports_item_interaction(Visual_Family family) noexcept { } bool uses_external_attributes(Visual_Family family) noexcept { - return family == Visual_Family::mesh || family == Visual_Family::marker; + switch (family) { + case Visual_Family::point: + case Visual_Family::splat: + case Visual_Family::pixel: + case Visual_Family::marker: + case Visual_Family::sphere: + case Visual_Family::primitive: + case Visual_Family::mesh: + case Visual_Family::path: + return true; + case Visual_Family::segment: + case Visual_Family::vector: + case Visual_Family::image: + case Visual_Family::labels: + case Visual_Family::glyph: + case Visual_Family::text: + case Visual_Family::volume: + return false; + } + return false; } template @@ -1354,9 +1373,23 @@ void Datoviz_Visual_Backend::ensure_external_attributes( throw std::length_error("Datoviz external attribute item count exceeds uint32"); std::size_t expected_count{}; - if (target.family == Visual_Family::mesh) expected_count = 3; - else if (target.family == Visual_Family::marker) expected_count = 5; - else throw std::logic_error("external attributes requested for an unsupported Visual family"); + switch (target.family) { + case Visual_Family::point: expected_count = 3; break; + case Visual_Family::splat: expected_count = 4; break; + case Visual_Family::pixel: expected_count = 3; break; + case Visual_Family::marker: expected_count = 5; break; + case Visual_Family::sphere: expected_count = 3; break; + case Visual_Family::segment: expected_count = 4; break; + case Visual_Family::vector: expected_count = 4; break; + case Visual_Family::primitive: + case Visual_Family::mesh: expected_count = 3; break; + case Visual_Family::path: expected_count = 3; break; + case Visual_Family::image: + case Visual_Family::labels: expected_count = 2; break; + default: + throw std::logic_error( + "external attributes requested for an unsupported Visual family"); + } if (target.attributes.size() == expected_count && std::ranges::all_of(target.attributes, [&](const External_Attribute& value) { return value.capacity >= item_count; @@ -1400,14 +1433,66 @@ void Datoviz_Visual_Backend::ensure_external_attributes( scene_buffer, gpu_buffer, false}); }; - create("position", sizeof(std::array)); - create("color", sizeof(std::array)); - if (target.family == Visual_Family::mesh) - create("normal", sizeof(std::array)); - else { + switch (target.family) { + case Visual_Family::point: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("diameter_px", sizeof(float)); + break; + case Visual_Family::splat: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("sigma", sizeof(std::array)); + create("angle", sizeof(float)); + break; + case Visual_Family::pixel: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("pixel_size_px", sizeof(float)); + break; + case Visual_Family::marker: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); create("diameter_px", sizeof(float)); create("angle", sizeof(float)); create("shape", sizeof(std::uint32_t)); + break; + case Visual_Family::sphere: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("radius", sizeof(float)); + break; + case Visual_Family::segment: + create("position_start", sizeof(Prepared_Position)); + create("position_end", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("stroke_width_px", sizeof(float)); + break; + case Visual_Family::vector: + create("position", sizeof(Prepared_Position)); + create("vector", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("stroke_width_px", sizeof(float)); + break; + case Visual_Family::primitive: + case Visual_Family::mesh: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("normal", sizeof(Prepared_Position)); + break; + case Visual_Family::path: + create("position", sizeof(Prepared_Position)); + create("color", sizeof(Prepared_Color)); + create("stroke_width_px", sizeof(float)); + break; + case Visual_Family::image: + case Visual_Family::labels: + create("position", sizeof(Prepared_Position)); + create("extent", sizeof(std::array)); + break; + default: + throw std::logic_error( + "external attribute layout requested for an unsupported Visual family"); } target.attributes = std::move(attributes); } @@ -1427,18 +1512,97 @@ void Datoviz_Visual_Backend::upload_external_attributes( dvz_buffer_upload(attribute.gpu_buffer, byte_offset, byte_count, source); }; auto attribute = target.attributes.begin(); - if (target.family == Visual_Family::mesh) { - const auto& data = prepared_data(visual); + switch (target.family) { + case Visual_Family::point: { + const auto& data = prepared_data(visual); upload(*attribute++, data.positions.data(), data.positions.size()); upload(*attribute++, data.colors.data(), data.colors.size()); - upload(*attribute++, data.normals.data(), data.normals.size()); - } else { + upload(*attribute++, data.diameters.data(), data.diameters.size()); + break; + } + case Visual_Family::splat: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.sigmas.data(), data.sigmas.size()); + upload(*attribute++, data.angles.data(), data.angles.size()); + break; + } + case Visual_Family::pixel: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.sizes.data(), data.sizes.size()); + break; + } + case Visual_Family::marker: { const auto& data = prepared_data(visual); upload(*attribute++, data.positions.data(), data.positions.size()); upload(*attribute++, data.colors.data(), data.colors.size()); upload(*attribute++, data.diameters.data(), data.diameters.size()); upload(*attribute++, data.angles.data(), data.angles.size()); upload(*attribute++, data.shapes.data(), data.shapes.size()); + break; + } + case Visual_Family::sphere: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.centers.data(), data.centers.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.radii.data(), data.radii.size()); + break; + } + case Visual_Family::segment: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.starts.data(), data.starts.size()); + upload(*attribute++, data.ends.data(), data.ends.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.widths.data(), data.widths.size()); + break; + } + case Visual_Family::vector: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.origins.data(), data.origins.size()); + upload(*attribute++, data.directions.data(), data.directions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.widths.data(), data.widths.size()); + break; + } + case Visual_Family::primitive: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.normals.data(), data.normals.size()); + break; + } + case Visual_Family::mesh: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.normals.data(), data.normals.size()); + break; + } + case Visual_Family::path: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.colors.data(), data.colors.size()); + upload(*attribute++, data.widths.data(), data.widths.size()); + break; + } + case Visual_Family::image: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.extents.data(), data.extents.size()); + break; + } + case Visual_Family::labels: { + const auto& data = prepared_data(visual); + upload(*attribute++, data.positions.data(), data.positions.size()); + upload(*attribute++, data.extents.data(), data.extents.size()); + break; + } + default: + throw std::logic_error( + "external attribute upload requested for an unsupported Visual family"); } } @@ -1566,13 +1730,19 @@ void Datoviz_Visual_Backend::apply_visual( } const auto count = static_cast(item_count); if (uses_external_attributes(target.family)) { + /* + * 直接顶点族在第一次录制前即绑定三槽外部属性。Segment/Vector/Image/ + * Labels 仍依赖 Datoviz 的 CPU 派生几何或纹理资源,不进入这条路径; + * 把 dense 数据先写入再改绑会违反 Datoviz 的单一属性所有权契约。 + */ ensure_external_attributes(target, point); if (bind_target) bind_external_attributes(target, target_index, count); upload_external_attributes(target, point, target_index); if (structure_changed && dvz_visual_set_visible(visual, point.visible) != DVZ_OK) - throw std::runtime_error("failed to apply Datoviz external visual visibility"); + throw std::runtime_error( + "failed to apply Datoviz external visual visibility"); target.applied_revision = point.revision; target.applied = point; return; diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp index b4d3af0..7b993b8 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp @@ -76,26 +76,22 @@ void Gpu_Completion_Service::cancel_reserved( std::lock_guard lock(pending->mutex); if (pending->status == Pending_Fence::Status::reserved) pending->status = Pending_Fence::Status::canceled; } -void Gpu_Completion_Service::acquire_slot() { - const auto started = std::chrono::steady_clock::now(); - if (!slots_.try_acquire()) { - backpressure_count_.fetch_add(1, std::memory_order_relaxed); - slots_.acquire(); - const auto waited = std::chrono::duration_cast( - std::chrono::steady_clock::now() - started) - .count(); - if (waited > 0) - backpressure_wait_ns_.fetch_add( - static_cast(waited), - std::memory_order_relaxed); +bool Gpu_Completion_Service::acquire_slot() noexcept { + std::size_t current = in_flight_.load(std::memory_order_relaxed); + while (current < static_cast(default_capacity)) { + if (in_flight_.compare_exchange_weak( + current, current + 1, + std::memory_order_acq_rel, + std::memory_order_relaxed)) { + update_peak(peak_in_flight_, current + 1); + return true; + } } - const std::size_t in_flight = - in_flight_.fetch_add(1, std::memory_order_relaxed) + 1; - update_peak(peak_in_flight_, in_flight); + backpressure_count_.fetch_add(1, std::memory_order_relaxed); + return false; } void Gpu_Completion_Service::release_slot() noexcept { in_flight_.fetch_sub(1, std::memory_order_relaxed); - slots_.release(); } Gpu_Completion_Service::Prepare_Result Gpu_Completion_Service::prepare( Completion completion, Exception_Handler on_exception, bool observe) { @@ -107,7 +103,10 @@ Gpu_Completion_Service::Prepare_Result Gpu_Completion_Service::prepare( pending->on_exception = std::move(on_exception); pending->observe = observe; pending->service = this; - acquire_slot(); + /* 唯一 GPU Submit 域绝不等待 completion 容量。满载时把背压作为 + * 准入结果立即反馈给 Scene,由上游下一帧策略自然重试。 */ + if (!acquire_slot()) + return {{}, Admission_Result::capacity_exhausted}; if (stopping_.load(std::memory_order_acquire)) { release_slot(); return {{}, Admission_Result::stopping}; diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp index 2553417..1f06e05 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp @@ -12,7 +12,6 @@ #include #include #include -#include #include namespace aethera::render_3d::detail { class Gpu_Completion_Service final { @@ -20,6 +19,7 @@ class Gpu_Completion_Service final { public: enum class Admission_Result : std::uint8_t { none, + capacity_exhausted, stopping }; enum class Completion_Error : std::uint8_t { @@ -84,7 +84,7 @@ private: static void update_max(std::atomic_uint64_t& maximum, std::uint64_t value) noexcept; static void cancel_reserved(const std::shared_ptr& pending); - void acquire_slot(); + [[nodiscard]] bool acquire_slot() noexcept; void release_slot() noexcept; void publish_state(std::size_t active_fences, std::size_t pending_fences) noexcept; @@ -94,7 +94,6 @@ private: static constexpr std::uint64_t fence_wait_timeout_ns = 1'000'000; static constexpr std::uint64_t maximum_fence_age_ns = 30'000'000'000ULL; static constexpr auto state_publication_interval = std::chrono::milliseconds(100); - std::counting_semaphore slots_{default_capacity}; /* 有界 reservation 槽位。 */ std::mutex pending_mutex_; /* 保护新登记 fence 队列。 */ std::deque> pending_; /* 完成线程尚未分组的 fence。 */ std::mutex wait_mutex_; /* 保护完成线程的条件等待。 */ diff --git a/render_3D/render_3D/scene/Render_Scene_3D.ipp b/render_3D/render_3D/scene/Render_Scene_3D.ipp index 4f63c80..d33f405 100644 --- a/render_3D/render_3D/scene/Render_Scene_3D.ipp +++ b/render_3D/render_3D/scene/Render_Scene_3D.ipp @@ -4,6 +4,9 @@ #include #include #include +#include +#include +#include namespace aethera::render_3d { namespace detail { template @@ -34,6 +37,9 @@ struct Render_Scene_3D::Private : Prev_Private { std::mutex completion_mutex{}; /* 只保护 completion 图拓扑的在途生命周期。 */ std::condition_variable completion_condition{}; bool completion_running{}; /* 析构等待异步 completion 图结束的唯一状态源。 */ + bool frame_topology_finished{}; + bool gpu_frame_finished{}; + bool frame_finish_started{}; std::shared_ptr backend{}; /* Scene 拥有的异步后端;已入队命令自行延长实现寿命。 */ std::shared_ptr paint_context{}; /* Paint 节点读取 Scene 当前 Prop 的生命周期门闩。 */ Root* camera_component{}; /* Builder 绑定的 Camera 组件。 */ @@ -42,6 +48,12 @@ struct Render_Scene_3D::Private : Prev_Private { std::array (*read_axes)(const Root*){}; /* 读取三轴当前配置。 */ Frame_3D* active_frame{}; /* 当前同步 process 借用的外部帧;提交完成后清空。 */ const Dispatch* dispatch{}; /* 最终 Scene 类型对应的静态公开分派表。 */ + std::unique_ptr submit_taskflow{}; + std::unique_ptr frame_taskflow{}; + Event_Batch active_events{}; + bool active_backend_submitted{}; + bool active_trace{}; + std::chrono::steady_clock::time_point active_prepare_started{}; ~Private(); /* Builder 内部初始化后端;必须在绑定 Visual Paint 目标之前调用一次。 */ template @@ -50,10 +62,10 @@ struct Render_Scene_3D::Private : Prev_Private { Camera_Descriptor (*camera_reader)(const Root*), std::array (*axes_reader)(const Root*)); template [[nodiscard]] detail::Scene_3D_Parameters parameters(Object* object) const; + template void complete_frame(Object* object, Frame_3D* frame); /* CRTP 覆盖:绑定 Scene 机制和 Render_Scene_3D 公开薄壳。 */ template void bind_private_crtp(Object* object); - /* CRTP 覆盖:先执行 Kernel Prepare Taskflow,再运行只负责异步入队的三维 Submit 阶段。 */ - template void process(Object* object, Callback&& callback) requires std::invocable; + template void ensure_frame_taskflow(Object* object); template [[nodiscard]] Render_Result render(Object* object, Frame_3D* frame); template [[nodiscard]] static const Dispatch& dispatch_for(); }; @@ -85,9 +97,11 @@ Render_Scene_3D::Builder::add_renderable(Visual_Object* visual_value) { visuals, visual_identity, &detail::Prepared_Visual_Instance::identity); if (found == visuals.end()) - visuals.push_back({visual_identity, std::move(erased)}); - else - found->visual = std::move(erased); + throw std::logic_error( + "3D Visual published outside its Scene registration"); + /* Builder 已经为每个 Visual 建立稳定槽位。并行 Submit 节点只写各自 + * 的 visual 成员,不扩容容器,也不互相读写同一对象。 */ + found->visual = std::move(erased); }); }; visuals.push_back(binding); @@ -158,6 +172,55 @@ detail::Scene_3D_Parameters Render_Scene_3D::Private::parameters(Object* object) return {prop.viewport, prop.clear_color, read_camera(camera_component), axis[0], axis[1], axis[2]}; } template +void Render_Scene_3D::Private::complete_frame(Object* object, Frame_3D* frame) { + Frame_Callback callback; + { + std::lock_guard lock(render_mutex); + callback = frame_callback; + } + auto finish = [this, object, frame, callback = std::move(callback)]() mutable { + aethera::detail::finish_taskflow_trace(*frame); + active_trace = false; + frame->mark(Frame_Trace_Marker::callback_started); + if (callback) callback(frame); + frame->mark(Frame_Trace_Marker::callback_finished); + frame->mark(Frame_Trace_Marker::frame_ready); + object->template publish_state( + [this, frame](State_Access states) { + auto& state = states.template get(); + state.frame_statistics = frame_statistics.submit( + *frame, Frame_Dimension::three_dimensional); + state.event_statistics = event_statistics.state(); + }); + auto context = std::static_pointer_cast< + detail::Scene_Paint_Context>(paint_context); + if (context->frame == frame) context->frame = nullptr; + if (active_frame == frame) active_frame = nullptr; + { + std::lock_guard lock(render_mutex); + frame_in_flight = false; + } + { + std::lock_guard lock(completion_mutex); + completion_running = false; + frame_topology_finished = false; + gpu_frame_finished = false; + frame_finish_started = false; + } + completion_condition.notify_all(); + }; + if (completion_graph.empty()) { + finish(); + return; + } + 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)); +} +template void Render_Scene_3D::Private::initialize_backend( Object* object, std::uint32_t gpu_index, bool validation_enabled, std::vector visuals, @@ -168,215 +231,367 @@ void Render_Scene_3D::Private::initialize_backend( read_camera = camera_reader; read_axes = axes_reader; const auto initial = parameters(object); + auto prepared_visuals = std::make_shared(); + prepared_visuals->reserve(visuals.size()); + for (const auto& visual : visuals) + prepared_visuals->push_back({visual.identity, {}}); backend = std::make_shared(gpu_index, validation_enabled, std::move(visuals), initial); backend->set_frame_callback([this, object](Frame_3D* frame) { - Frame_Callback callback; - { - std::lock_guard lock(render_mutex); - callback = frame_callback; - } - auto finish = [this, object, frame, callback = std::move(callback)]() mutable { - aethera::detail::finish_taskflow_trace(*frame); - frame->mark(Frame_Trace_Marker::callback_started); - if (callback) callback(frame); - frame->mark(Frame_Trace_Marker::callback_finished); - frame->mark(Frame_Trace_Marker::frame_ready); - object->template publish_state( - [this, frame](State_Access states) { - auto& state = states.template get(); - state.frame_statistics = frame_statistics.submit( - *frame, Frame_Dimension::three_dimensional); - state.event_statistics = event_statistics.state(); - }); - { - std::lock_guard lock(render_mutex); - frame_in_flight = false; - } - { - std::lock_guard lock(completion_mutex); - completion_running = false; - } - completion_condition.notify_all(); - }; - if (completion_graph.empty()) { - finish(); - return; - } + bool finish{}; { std::lock_guard lock(completion_mutex); - completion_running = true; + gpu_frame_finished = true; + if (frame_topology_finished && !frame_finish_started) { + frame_finish_started = true; + finish = true; + } } - /* - * GPU 完成线程只触发 Taskflow topology,不等待编码/发送节点;最终回调在 topology 完成后执行。 - * 帧准入在 finish() 之前始终有效,因此 completion 图运行期间该外部 Frame 不会被下一帧复用。 - */ - 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)); + if (finish) complete_frame(object, frame); }); paint_context = std::make_shared>( detail::Scene_Paint_Context{ backend, object, nullptr, initial, - std::make_shared()}); + std::move(prepared_visuals)}); } -template -void Render_Scene_3D::Private::process(Object* object, Callback&& callback) requires std::invocable { - Frame_3D* frame = active_frame; - if (!frame) throw std::logic_error("3D Scene process has no external frame"); - const auto& prop = object->template read_prop(); - if (!prop.view_active || prop.viewport.empty() || !backend || !backend->available()) return; - frame->mark(Frame_Trace_Marker::scene_render_started); - frame->mark(Frame_Trace_Marker::event_dispatch_started); - auto events = this->take_events(object); - frame->mark(Frame_Trace_Marker::event_dispatch_finished); - auto context = std::static_pointer_cast>(paint_context); - context->parameters = parameters(object); - /* - * Visual 的 Prepared 快照只由数据变化更新,并持久保存在 Scene_Paint_Context。 - * 静态帧不推进/复制 Renderable 双缓冲,也不启动 Kernel Taskflow,只提交已经准备好的 - * 完整快照。动态 Visual 变脏时才进入完整 Prepare/Submit 路径并原位更新这份缓存。 - */ - bool prepare_dirty = context->visuals->empty(); - const auto current_prepare = object->template current_dependency_graph(); - current_prepare.for_each_bound([&](Renderable* renderable, Renderable::Private&) { - prepare_dirty = prepare_dirty || renderable->template dirty(); - }); - bool submit_dirty{}; - const auto current_submit = object->template current_dependency_graph(); - current_submit.for_each_bound([&](Renderable* renderable, Renderable::Private&) { - submit_dirty = submit_dirty || renderable->template dirty(); - }); - if (!prepare_dirty && !submit_dirty) { - frame->mark(Frame_Trace_Marker::prepare_started); - frame->mark(Frame_Trace_Marker::prepare_finished); - frame->mark(Frame_Trace_Marker::paint_started); - const auto result = context->backend->render( - context->visuals, context->parameters, frame, std::move(events)); - frame->mark(Frame_Trace_Marker::paint_finished); - if (result == detail::Async_Render_Backend::Submit_Result::queued) - std::invoke(std::forward(callback)); - return; - } - - /* 后端仍可能持有上一版快照时只在真正变脏的帧执行一次 COW。 - * 静态帧仅复制 shared_ptr,禁止把完整 Prepared Visual 集合复制进提交队列。 */ - if (context->visuals.use_count() != 1) - context->visuals = - std::make_shared(*context->visuals); - camera_component->advance_object(); - axes_component->advance_object(); - struct Frame_Context_Scope { - Frame_3D*& target; /* Visual Submit 回调读取的当前外部帧槽位。 */ - Frame_3D* previous{}; /* 嵌套调用前的帧;析构时恢复。 */ - ~Frame_Context_Scope() { target = previous; } - } frame_context_scope{context->frame, context->frame}; - context->frame = frame; - bool backend_submitted{}; - Prev_Private::process(object, frame, [&](const Scene::Private::Result&) { - frame->mark(Frame_Trace_Marker::paint_started); - const auto prepare = object->template current_dependency_graph(); - prepare.for_each_bound([](Renderable* renderable, Renderable::Private& data) { if (data.dispatch->state.pending(renderable)->prepare_executed) renderable->template mark_dirty(); }); - const auto submit = object->template current_dependency_graph(); - const auto result = submit.for_each_topological_view([&](const auto& view, const Dependency_Graph::Node& node) { +template +void Render_Scene_3D::Private::ensure_frame_taskflow(Object* object) { + if (frame_taskflow) return; + if (!runtime->taskflow) + runtime->taskflow = std::make_unique("scene.prepare"); + /* 模块图在外层 topology 提交前必须已经包含真实节点。若把首次 builder + * 留给模块内部 condition,Taskflow 已经物化的父 topology 看不到本次新增 + * 节点,Visual 会发布尚未 Prepare 的空 payload。 */ + object->template current_dependency_graph().for_each_bound( + [](Renderable* renderable, Renderable::Private& data) { + auto* dispatch = data.dispatch; + if (!dispatch->prepare.builder || data.prepare_graph_built) return; + if (!data.prepare_graph) + data.prepare_graph = std::make_unique( + std::string(dispatch->business_name) + ".prepare.graph"); + *data.prepare_graph = dispatch->prepare.builder(renderable); + data.prepare_graph_built = true; + renderable->template mark_dirty(); + }); + submit_taskflow = std::make_unique("render_3d.visual.submit"); + struct Submit_Tasks { + Task_Node entry; + Task_Node exit; + }; + std::unordered_map submit_tasks; + const auto submit_dependencies = object->template current_dependency_graph(); + const auto order = submit_dependencies.for_each_topological_view( + [&](const auto& view, const Dependency_Graph::Node& node) { auto* data = view.private_data(node); if (!data) return; Root* root = node.object; auto* dispatch = data->dispatch; - auto& state = *dispatch->state.pending(root); - state.paint_graph_rebuilt = false; - state.paint_execution_time_ns = 0; - const bool dirty = root->template dirty(); - state.paint_dirty = dirty; - state.paint_executed = dispatch->paint.predicate(root, dirty); - 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( - 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->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); + const auto prefix = std::string(dispatch->business_name) + ".submit"; + if (dispatch->paint.builder && !data->paint_graph) { + data->paint_graph = std::make_unique(prefix + ".graph"); + *data->paint_graph = dispatch->paint.builder(root); + data->paint_graph_built = true; + } + auto condition = submit_taskflow->add_condition(prefix + ".condition", + [data, dispatch, root] { + auto& state = *dispatch->state.pending(root); + state.paint_graph_rebuilt = false; + state.paint_execution_time_ns = 0; + const bool dirty = root->template dirty(); + state.paint_dirty = dirty; + state.paint_executed = dispatch->paint.predicate(root, dirty); + if (state.paint_executed) + state.paint_execution_time_ns = static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()).count()); + return state.paint_executed ? 0 : 1; + }); + condition.describe("renderable", std::string(dispatch->business_name)) + .describe("dimension", "3D") + .describe("stage", "visual submit predicate"); + Task_Node publish; + if (dispatch->paint.builder) + publish = submit_taskflow->compose(prefix + ".visual", *data->paint_graph); + else + publish = submit_taskflow->add(prefix + ".visual", [dispatch, root] { + if (dispatch->paint.run) dispatch->paint.run(root); + }); + publish.describe("renderable", std::string(dispatch->business_name)) + .describe("dimension", "3D") + .describe("stage", "prepared visual publish"); + auto done = submit_taskflow->add(prefix + ".complete", [dispatch, root] { + auto& state = *dispatch->state.pending(root); + if (state.paint_executed) { + root->template take_dirty(); + const auto finished = static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()).count()); + state.paint_execution_time_ns = finished - state.paint_execution_time_ns; } - } else { state.paint_task_count = dispatch->paint.run ? 1 : 0; if (dispatch->paint.run) dispatch->paint.run(root); } - if (!data->paint_extension.empty()) - 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( - std::chrono::steady_clock::now() - paint_started).count()); - dispatch->state.publish(root); + dispatch->state.publish(root); + }); + done.describe("renderable", std::string(dispatch->business_name)) + .describe("dimension", "3D") + .describe("stage", "visual state publish"); + condition.precede(publish); + condition.precede(done); + if (data->paint_extension.empty()) + publish.precede(done); + else { + auto extension = submit_taskflow->compose( + prefix + ".extension", data->paint_extension); + extension.describe("owner", std::string(dispatch->business_name)) + .describe("stage", "renderable extension"); + publish.precede(extension); + extension.precede(done); + } + submit_tasks.emplace(root, Submit_Tasks{condition, done}); }); - if (!result) throw std::logic_error("render scene submit graph became invalid during submission"); - if (context->visuals->empty()) - throw std::logic_error("render scene submit graph published no Visual snapshots"); - backend_submitted = context->backend->render( - context->visuals, context->parameters, context->frame, - std::move(events)) == - detail::Async_Render_Backend::Submit_Result::queued; - frame->mark(Frame_Trace_Marker::paint_finished); + if (!order) + throw std::logic_error("render scene submit graph became invalid while building"); + /* Submit 只继承业务依赖图中的真实前驱;互不依赖的 Visual 不再被人为 + * 串成一个长条。未绑定的中间依赖节点只用于传递依赖关系。 */ + submit_dependencies.for_each( + [&](const Dependency_Graph::Node& target_node) { + if (!submit_dependencies.private_data(target_node)) return; + const auto target = submit_tasks.find(target_node.object); + if (target == submit_tasks.end()) return; + std::unordered_set visited; + std::vector pending{ + target_node.dependencies.begin(), target_node.dependencies.end()}; + while (!pending.empty()) { + const auto* dependency = pending.back(); + pending.pop_back(); + if (!visited.insert(dependency->object).second) continue; + if (submit_dependencies.private_data(*dependency)) { + const auto source = submit_tasks.find(dependency->object); + if (source != submit_tasks.end()) + source->second.exit.precede(target->second.entry); + continue; + } + pending.insert(pending.end(), dependency->dependencies.begin(), + dependency->dependencies.end()); + } + }); + + frame_taskflow = std::make_unique("render_3d.frame"); + auto& graph = *frame_taskflow; + auto begin = graph.add("scene.begin", [this, object] { + if (!active_frame) throw std::logic_error("3D frame DAG lost its active frame"); + active_frame->mark(Frame_Trace_Marker::scene_render_started); + active_frame->mark(Frame_Trace_Marker::event_dispatch_started); + active_events = take_events(object); + active_frame->mark(Frame_Trace_Marker::event_dispatch_finished); + auto context = std::static_pointer_cast>(paint_context); + context->parameters = parameters(object); + if (!detail::prepared_visual_batch_complete(*context->visuals)) + object->template current_dependency_graph().for_each_bound( + [](Renderable* renderable, Renderable::Private&) { + renderable->template mark_dirty(); + }); + if (context->visuals.use_count() != 1) + context->visuals = std::make_shared(*context->visuals); + context->frame = active_frame; + camera_component->advance_object(); + axes_component->advance_object(); + active_prepare_started = std::chrono::steady_clock::now(); + active_frame->mark(Frame_Trace_Marker::prepare_started); }); - if (backend_submitted) std::invoke(std::forward(callback)); + begin.describe("dimension", "3D").describe("stage", "event and frame setup"); + struct Prepare_Tasks { + Task_Node task; + }; + std::unordered_map prepare_tasks; + const auto prepare_dependencies = + object->template current_dependency_graph(); + const auto prepare_order = prepare_dependencies.for_each_topological_view( + [&](const auto& view, const Dependency_Graph::Node& node) { + auto* data = view.private_data(node); + if (!data) return; + Root* root = node.object; + auto* dispatch = data->dispatch; + auto task = graph.add( + std::string(dispatch->business_name) + ".prepare.data", + [dispatch, root] { + auto& state = *dispatch->state.pending(root); + const bool dirty = root->template dirty(); + state.prepare_dirty = dirty; + state.prepare_graph_rebuilt = false; + state.prepare_task_count = dispatch->prepare.run ? 1 : 0; + state.prepare_executed = + dispatch->prepare.predicate(root, dirty); + state.prepare_execution_time_ns = 0; + if (state.prepare_executed) { + const auto started = std::chrono::steady_clock::now(); + if (dispatch->prepare.run) dispatch->prepare.run(root); + root->template take_dirty(); + state.prepare_execution_time_ns = + static_cast( + std::chrono::duration_cast< + std::chrono::nanoseconds>( + std::chrono::steady_clock::now() - started) + .count()); + } + dispatch->state.publish(root); + }); + task.describe("renderable", std::string(dispatch->business_name)) + .describe("dimension", "3D") + .describe("stage", "CPU prepared data") + .describe("execution_domain", "Taskflow worker"); + prepare_tasks.emplace(root, Prepare_Tasks{task}); + }); + if (!prepare_order) + throw std::logic_error( + "render scene prepare graph became invalid while building"); + prepare_dependencies.for_each( + [&](const Dependency_Graph::Node& target_node) { + if (!prepare_dependencies.private_data(target_node)) return; + const auto target = prepare_tasks.find(target_node.object); + if (target == prepare_tasks.end()) return; + std::unordered_set visited; + std::vector pending{ + target_node.dependencies.begin(), target_node.dependencies.end()}; + while (!pending.empty()) { + const auto* dependency = pending.back(); + pending.pop_back(); + if (!visited.insert(dependency->object).second) continue; + if (prepare_dependencies.private_data(*dependency)) { + const auto source = prepare_tasks.find(dependency->object); + if (source != prepare_tasks.end()) + source->second.task.precede(target->second.task); + continue; + } + pending.insert(pending.end(), dependency->dependencies.begin(), + dependency->dependencies.end()); + } + }); + for (const auto& [root, tasks] : prepare_tasks) { + static_cast(root); + begin.precede(tasks.task); + } + auto prepare_done = graph.add("scene.prepare.complete", [this, object] { + active_frame->mark(Frame_Trace_Marker::prepare_finished); + const auto context = std::static_pointer_cast< + detail::Scene_Paint_Context>(paint_context); + const bool initial_publish = + !detail::prepared_visual_batch_complete(*context->visuals); + const auto prepare_graph = object->template current_dependency_graph(); + prepare_graph.for_each_bound([initial_publish](Renderable* renderable, Renderable::Private& data) { + if (initial_publish || data.dispatch->state.pending(renderable)->prepare_executed) + renderable->template mark_dirty(); + }); + auto& private_data = static_cast(*this); + auto& state = static_cast(*private_data.state.pending); + state.taskflow_execution_time_ns = static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now() - active_prepare_started).count()); + state.event_statistics = event_statistics.state(); + active_frame->mark(Frame_Trace_Marker::paint_started); + }); + prepare_done.describe("dimension", "3D") + .describe("stage", "prepare state publish") + .describe("state_direction", "internal to external"); + auto visuals = graph.compose("scene.visuals", *submit_taskflow); + visuals.describe("dimension", "3D") + .describe("stage", "prepared visual collection") + .describe("execution_domain", "Taskflow workers") + .describe("node_count", std::to_string(submit_taskflow->size())); + auto submit = graph.add("scene.backend.submit", [this, object] { + auto context = std::static_pointer_cast>(paint_context); + if (!detail::prepared_visual_batch_complete(*context->visuals)) + throw std::logic_error( + "render scene published an incomplete Visual batch"); + active_backend_submitted = context->backend->render( + context->visuals, context->parameters, active_frame, + std::move(active_events)) == + detail::Async_Render_Backend::Submit_Result::queued; + active_frame->mark(Frame_Trace_Marker::paint_finished); + if (active_backend_submitted) + active_frame->mark(Frame_Trace_Marker::scene_render_finished); + }); + submit.describe("dimension", "3D") + .describe("backend", "Datoviz") + .describe("stage", "backend prepare and submission enqueue") + .describe("cpu_owner", "Taskflow worker") + .describe("gpu_submit_owner", "GPU Render Domain") + .describe("completion", "GPU fence callback") + .describe("output", "RGBA8 readback"); + for (const auto& [root, tasks] : prepare_tasks) { + static_cast(root); + tasks.task.precede(prepare_done); + } + prepare_done.precede(visuals); + visuals.precede(submit); } + template Render_Scene_3D::Render_Result Render_Scene_3D::Private::render(Object* object, Frame_3D* frame) { if (!frame) throw std::invalid_argument("Render_Scene_3D requires a non-null external frame"); + 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; + if (!backend || !backend->available()) return Render_Result::backend_unavailable; + ensure_frame_taskflow(object); { std::lock_guard lock(render_mutex); if (!frame_callback) throw std::logic_error("Render_Scene_3D requires a frame callback before render"); if (frame_in_flight) return Render_Result::frame_in_flight; frame_in_flight = true; + active_frame = frame; + active_backend_submitted = false; } - bool submitted{}; - struct Admission_Scope { - Private& data; - bool& submitted; - ~Admission_Scope() { - if (submitted) return; - std::lock_guard lock(data.render_mutex); - data.frame_in_flight = false; - } - } 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{}; /* 嵌套调用前的帧;析构时恢复。 */ - ~Active_Frame_Scope() { target = previous; } - } active_frame_scope{active_frame, active_frame}; - active_frame = frame; - object->process([&] { submitted = true; frame->mark(Frame_Trace_Marker::scene_render_finished); }); - if (submitted) { - trace_scope.active = false; - return Render_Result::submitted; + active_trace = aethera::detail::begin_taskflow_trace(*frame); + { + std::lock_guard lock(completion_mutex); + completion_running = true; + frame_topology_finished = false; + gpu_frame_finished = false; + frame_finish_started = false; } - 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; - return Render_Result::backend_unavailable; + try { + auto completion = [this, object, frame] { + auto& private_data = static_cast(*this); + private_data.state.advance(); + object->template notify_state(); + bool finish{}; + { + std::lock_guard lock(completion_mutex); + frame_topology_finished = true; + if ((!active_backend_submitted || gpu_frame_finished) && + !frame_finish_started) { + frame_finish_started = true; + finish = true; + } + } + if (finish) complete_frame(object, frame); + }; + if (frame->taskflow_trace_requested()) + aethera::detail::run_taskflow( + *frame_taskflow, *frame, "render_3d.frame", std::move(completion)); + else + aethera::detail::run_taskflow(*frame_taskflow, std::move(completion)); + } + catch (...) { + if (active_trace) { + aethera::detail::finish_taskflow_trace(*frame); + active_trace = false; + } + active_frame = nullptr; + std::lock_guard lock(render_mutex); + frame_in_flight = false; + { + std::lock_guard completion_lock(completion_mutex); + completion_running = false; + frame_topology_finished = false; + gpu_frame_finished = false; + frame_finish_started = false; + } + completion_condition.notify_all(); + throw; + } + return Render_Result::submitted; } template const Render_Scene_3D::Private::Dispatch& Render_Scene_3D::Private::dispatch_for() { static const Dispatch value{[](Root* root) { auto* object = static_cast(root); auto& data = static_cast(*object->d); object->template publish_state([&data](State_Access states) { data.frame_statistics.reset(); data.event_statistics.reset(); auto& state = states.template get(); state.frame_statistics = {}; state.event_statistics = {}; }); }, [](Root* root, Frame_3D* frame) { auto* object = static_cast(root); return static_cast(*object->d).render(object, frame); }, [](Root* root, Frame_Callback callback) { auto* object = static_cast(root); auto& data = static_cast(*object->d); std::lock_guard lock(data.render_mutex); data.frame_callback = std::move(callback); }, [](Root* root, bool active) { static_cast(root)->template set<&Prop::view_active>(active); }}; return value; } template void Render_Scene_3D::Private::bind_private_crtp(Object* object) { Prev_Private::bind_private_crtp(object); dispatch = &dispatch_for(); } diff --git a/web_server/src/Gallery_WebSocket.cpp b/web_server/src/Gallery_WebSocket.cpp index effdc36..0e98d85 100644 --- a/web_server/src/Gallery_WebSocket.cpp +++ b/web_server/src/Gallery_WebSocket.cpp @@ -153,9 +153,14 @@ void Gallery_WebSocket::close() noexcept { } catch (...) {} d->subscription = 0; - if (d->video) d->video->close(); + /* + * The exclusive lease describes live browser connections, not WebRTC + * teardown progress. Release it at the connection-close boundary so a + * closed page cannot block the next page while a media sender is retiring. + */ try { release_exclusive_page(d->page_id); } catch (...) {} + if (d->video) d->video->close(); } Gallery_WebSocket_Controller::Gallery_WebSocket_Controller( diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index 0cefff5..5a3789c 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -203,6 +203,10 @@ void append_event_statistics_json(nlohmann::json& output, } nlohmann::json taskflow_trace_json(const Taskflow_Frame_Trace& trace) { + nlohmann::json markers = nlohmann::json::object(); + for (const auto& marker : trace.markers) + markers[magic_enum::enum_name(marker.marker)] = + static_cast(marker.elapsed_ns) / 1'000'000.0; nlohmann::json graphs = nlohmann::json::array(); std::unordered_map node_ids; for (const auto& graph : trace.graphs) { @@ -215,11 +219,15 @@ nlohmann::json taskflow_trace_json(const Taskflow_Frame_Trace& trace) { nlohmann::json successors = nlohmann::json::array(); for (const auto native_id : node.successors) successors.push_back(std::to_string(native_id)); + nlohmann::json attributes = nlohmann::json::object(); + for (const auto& [key, value] : node.attributes) + attributes[key] = value; 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)}}); + {"successors", std::move(successors)}, + {"attributes", std::move(attributes)}}); } graphs.push_back({ {"stage", graph.stage}, {"name", graph.taskflow_name}, @@ -236,9 +244,15 @@ nlohmann::json taskflow_trace_json(const Taskflow_Frame_Trace& trace) { {"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}, + {"ready_ms", task.ready_ms}, {"entered_ms", task.entered_ms}, + {"started_ms", task.started_ms}, {"finished_ms", task.finished_ms}, + {"completed_ms", task.completed_ms}, {"duration_ms", task.duration_ms}, + {"cpu_duration_ms", task.cpu_duration_ms}, + {"observer_entry_ms", task.observer_entry_ms}, + {"observer_exit_ms", task.observer_exit_ms}, + {"observer_entry_cpu_ms", task.observer_entry_cpu_ms}, + {"observer_exit_cpu_ms", task.observer_exit_cpu_ms}, {"queue_wait_ms", task.queue_wait_ms}}); } return { @@ -246,6 +260,7 @@ nlohmann::json taskflow_trace_json(const Taskflow_Frame_Trace& trace) { {"correlation_id", trace.identity.correlation_id}, {"created_time_unix_ns", trace.created_time_unix_ns}, {"worker_count", trace.worker_count}, + {"markers", std::move(markers)}, {"graphs", std::move(graphs)}, {"executions", std::move(executions)}}; } @@ -344,7 +359,6 @@ struct Plot::Private { mutable std::mutex tick_mutex; std::optional pending_tick{}; /* 时钟拥塞时只保留尚未处理的最新时间点。 */ bool tick_task_scheduled{}; /* Taskflow 中是否已有唯一 tick 消费任务。 */ - 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 总数。 */ @@ -506,24 +520,56 @@ bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) { } 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()) { - 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(); if (!pacing.render_enabled || streams.consumers.empty()) return; + + std::size_t slot_index{}; + Managed_Frame* managed{}; { std::lock_guard lock(frame_mutex); - if (std::ranges::none_of(frame_slots, [](const Managed_Frame& slot) { - return slot.state == Frame_State::available; + /* + * Frame_State 是 Plot 数据准备生命周期的唯一权威来源。必须在修改 + * Scene/Visual 输入之前取得准入;Scene::render() 内部再拒绝已经太晚, + * 因为上一帧的异步 prepare 可能正在读取同一份业务数据。 + */ + if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) { + return slot.state == Frame_State::in_flight; })) { + preparation_busy_count.fetch_add(1, std::memory_order_relaxed); + return; + } + const auto available = std::ranges::find_if( + frame_slots, [](const Managed_Frame& slot) { + return slot.state == Frame_State::available; + }); + if (available == frame_slots.end()) { frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed); return; } + slot_index = static_cast( + std::distance(frame_slots.begin(), available)); + managed = &*available; + managed->state = Frame_State::in_flight; + managed->presentation_time = + std::chrono::duration_cast( + std::chrono::duration( + tick.time_milliseconds)); } + const auto rollback_unsubmitted = [this, slot_index] { + std::lock_guard lock(frame_mutex); + auto& slot = frame_slots[slot_index]; + 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 { tick.width = streams.width; tick.height = streams.height; /* @@ -550,37 +596,6 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { const std::uint64_t sequence = next_frame_sequence++; const Frame_Identity identity{sequence, tick.sequence == 0 ? sequence : tick.sequence}; - std::size_t slot_index{}; - Managed_Frame* managed{}; - { - std::lock_guard lock(frame_mutex); - const auto available = std::ranges::find_if( - frame_slots, [](const Managed_Frame& slot) { - return slot.state == Frame_State::available; - }); - if (available == frame_slots.end()) return; - slot_index = static_cast( - std::distance(frame_slots.begin(), available)); - managed = &*available; - managed->state = Frame_State::in_flight; - managed->presentation_time = - std::chrono::duration_cast( - std::chrono::duration( - tick.time_milliseconds)); - } - const auto rollback_unsubmitted = [this, slot_index] { - std::lock_guard lock(frame_mutex); - auto& slot = frame_slots[slot_index]; - 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)) { auto& output = *std::get>(managed->frame); output.begin(identity, pacing.video_enabled @@ -1007,6 +1022,8 @@ nlohmann::json Plot::diagnostics() const { {"abandoned_count", gpu.abandoned_count}, {"stopping", gpu.stopping}}; } + if (const auto failure = d->terminal_failure.load(std::memory_order_acquire)) + output["terminal_failure"] = *failure; return output; } diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index 515fd97..5c9ad63 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -74,13 +74,15 @@ type Frame_Metrics = {sequence: number; generated_time_unix_ms: number; server_c 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[]}; + predecessors: string[]; successors: string[]; attributes?: Record}; 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}; + worker_queue_capacity: number; ready_ms: number; entered_ms: number; started_ms: number; finished_ms: number; + completed_ms: number; duration_ms: number; cpu_duration_ms: number; observer_entry_ms: number; observer_exit_ms: number; + observer_entry_cpu_ms: number; observer_exit_cpu_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[]}; + markers?: Record; 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; @@ -858,7 +860,11 @@ type Pipeline_Stage_Definition = [string, string, string]; const pipeline_2d_definitions: Pipeline_Stage_Definition[] = [ ["pipeline_2d_event_ms", "2D 事件分发", "Scene 将当前输入事件分发给二维 Renderable 的耗时。"], ["pipeline_2d_prepare_ms", "2D 数据准备", "二维 Prepare 依赖图更新缓存、坐标映射和绘制数据的耗时。"], - ["pipeline_2d_paint_ms", "Blend2D 绘制", "二维 Paint 任务图清屏并写入 Blend2D 帧缓存的耗时。"], + ["pipeline_2d_frame_target_ms", "2D 帧目标准备", "调整主帧尺寸并清空完整 Blend2D 像素目标的耗时。"], + ["pipeline_2d_background_ms", "2D 背景填充", "向主帧目标填充 Scene 背景色的耗时。"], + ["pipeline_2d_cache_targets_ms", "2D 缓存目标准备", "检查缓存组并清空本帧需要重建的 Renderable 缓存目标。"], + ["pipeline_2d_taskflow_ms", "2D Taskflow 墙钟", "从提交 2D Paint Taskflow 到同步完成的墙钟时间,包含节点执行和真实 Executor 排队。"], + ["pipeline_2d_paint_coordination_ms", "2D Paint 阶段衔接", "Paint 总区间中不属于帧目标、背景、缓存目标和 Taskflow 的轻量衔接。"], ["pipeline_2d_scene_coordination_ms", "2D Scene 编排", "Scene 渲染区间内除事件、Prepare、Paint 外的依赖图编排耗时。"], ["pipeline_2d_callback_ms", "2D 完成帧提取", "同步二维帧完成回调中取得连续 RGBA8 像素并建立不可变完成帧的耗时。"], ["pipeline_2d_frame_handoff_ms", "2D 帧建立与完成发布", "Render_Frame 建立、Scene 入口以及完成帧交给页面图集之间尚未由独立 marker 覆盖的衔接耗时。"] @@ -921,7 +927,9 @@ function Frame_Timeline_Chart({diagnostics, dimension, paused, on_context_menu}: const series_keys: Array<[string, string, string]> = dimension === "2D" ? [ total_series, ["pipeline_2d_prepare_ms", "2D Prepare", "#62a8ff"], - ["pipeline_2d_paint_ms", "Blend2D 绘制", "#f4bd63"], + ["pipeline_2d_frame_target_ms", "帧目标准备", "#79d5ff"], + ["pipeline_2d_cache_targets_ms", "缓存目标准备", "#f4bd63"], + ["pipeline_2d_taskflow_ms", "Paint Taskflow", "#ef9f55"], ["pipeline_2d_scene_coordination_ms", "Scene 编排", "#ff7d9c"], ["pipeline_2d_frame_handoff_ms", "完成发布衔接", "#b998ff"] ] : [ @@ -1182,28 +1190,126 @@ function nanoseconds(value: number) { return milliseconds(value / 1_000_000); } -function Taskflow_Dag({graph, executions}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]}) { +const taskflow_operation_names: Record = { + condition: "执行条件", + data: "绘制", + graph: "子图", + extension: "扩展", + complete: "提交组件 State", + cache_composite: "缓存合成" +}; + +function taskflow_node_name(name: string) { + const parts = name.split(".").filter(part => part && part !== "paint"); + if (parts.length === 0) return name; + const operation = taskflow_operation_names[parts.at(-1) ?? ""] ?? parts.at(-1); + return parts.length === 1 ? operation ?? name : `${parts[0]} · ${operation}`; +} + +function taskflow_graph_analysis(graph: Taskflow_Graph_Trace, executions: Taskflow_Execution_Trace[]) { + const native_ids = new Set(graph.nodes.map(node => node.native_id)); + const rows = executions.filter(value => native_ids.has(value.native_id)); + const wall_time = Math.max(0, graph.finished_ms - graph.submitted_ms); + const execution_time = rows.reduce((sum, row) => sum + row.duration_ms, 0); + const cpu_execution_time = rows.reduce((sum, row) => sum + row.cpu_duration_ms, 0); + const descheduled_time = Math.max(0, execution_time - cpu_execution_time); + const observer_entry_time = rows.reduce((sum, row) => sum + row.observer_entry_ms, 0); + const observer_exit_time = rows.reduce((sum, row) => sum + row.observer_exit_ms, 0); + const observer_cpu_time = rows.reduce((sum, row) => sum + row.observer_entry_cpu_ms + row.observer_exit_cpu_ms, 0); + const queue_time = rows.reduce((sum, row) => sum + row.queue_wait_ms, 0); + const first_entered = rows.length ? Math.min(...rows.map(row => row.entered_ms)) : graph.submitted_ms; + const last_completed = rows.length ? Math.max(...rows.map(row => row.completed_ms)) : graph.finished_ms; + const longest = rows.reduce( + (result, row) => !result || row.duration_ms > result.duration_ms ? row : result, null); + const events = rows.flatMap(row => [ + {time: row.started_ms, delta: 1}, + {time: row.finished_ms, delta: -1} + ]).sort((left, right) => left.time - right.time || left.delta - right.delta); + let active = 0; + let maximum_parallelism = 0; + for (const event of events) { + active += event.delta; + maximum_parallelism = Math.max(maximum_parallelism, active); + } + const boundaries = [graph.submitted_ms, graph.finished_ms, + ...rows.flatMap(row => [row.entered_ms, row.started_ms, row.finished_ms, row.completed_ms])] + .filter(value => value >= graph.submitted_ms && value <= graph.finished_ms) + .sort((left, right) => left - right); + let body_wall_time = 0; + let observer_wall_time = 0; + let idle_wall_time = 0; + for (let index = 1; index < boundaries.length; ++index) { + const first = boundaries[index - 1]; + const last = boundaries[index]; + if (last <= first) continue; + const middle = (first + last) * .5; + if (rows.some(row => row.started_ms <= middle && middle < row.finished_ms)) + body_wall_time += last - first; + else if (rows.some(row => (row.entered_ms <= middle && middle < row.started_ms) || + (row.finished_ms <= middle && middle < row.completed_ms))) + observer_wall_time += last - first; + else idle_wall_time += last - first; + } + const initial_wait = Math.max(0, first_entered - graph.submitted_ms); + const completion_tail = Math.max(0, graph.finished_ms - last_completed); + const node_by_id = new Map(graph.nodes.map(node => [node.native_id, node])); + const levels = new Map(); + const visiting = new Set(); + const level_of = (native_id: string): number => { + const known = levels.get(native_id); + if (known !== undefined) return known; + if (visiting.has(native_id)) return 0; + visiting.add(native_id); + const node = node_by_id.get(native_id); + const level = node?.predecessors.length + ? 1 + Math.max(...node.predecessors.map(predecessor => level_of(predecessor))) + : 0; + visiting.delete(native_id); + levels.set(native_id, level); + return level; + }; + for (const node of graph.nodes) level_of(node.native_id); + const width_by_level = new Map(); + for (const level of levels.values()) width_by_level.set(level, (width_by_level.get(level) ?? 0) + 1); + return { + rows, levels, wall_time, execution_time, cpu_execution_time, + descheduled_time, observer_entry_time, observer_cpu_time, + observer_exit_time, queue_time, body_wall_time, observer_wall_time, + idle_wall_time, initial_wait, completion_tail, + internal_idle_time: Math.max(0, idle_wall_time - initial_wait - completion_tail), + peak_entry_queue: rows.reduce((maximum, row) => Math.max(maximum, row.worker_queue_size), 0), + maximum_parallelism, + layer_count: width_by_level.size, + parallel_layer_count: [...width_by_level.values()].filter(width => width > 1).length, + longest + }; +} + +function Taskflow_Dag({graph, executions, frame}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]; frame: Taskflow_Frame_Trace}) { const [nodes, set_nodes] = useState([]); const [edges, set_edges] = useState([]); + const [copy_state, set_copy_state] = useState("复制拓扑 JSON"); 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 analysis = taskflow_graph_analysis(graph, executions); 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, + flow_edges.push({id: `${node.id}->${target}`, source: node.id, target, type: "smoothstep", 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" + "elk.algorithm": "layered", "elk.direction": "DOWN", + "elk.edgeRouting": "ORTHOGONAL", + "elk.layered.spacing.nodeNodeBetweenLayers": "68", "elk.spacing.nodeNode": "42", + "elk.layered.nodePlacement.strategy": "BRANDES_KOEPF" }, - children: graph.nodes.map(node => ({id: node.id, width: 230, height: 92})), + children: graph.nodes.map(node => ({id: node.id, width: 280, height: 132})), edges: flow_edges.map(edge => ({id: edge.id, sources: [edge.source], targets: [edge.target]})) }).then(layout => { if (cancelled) return; @@ -1212,12 +1318,19 @@ function Taskflow_Dag({graph, executions}: {graph: Taskflow_Graph_Trace; executi const sample = execution.get(node.native_id); const wait = sample?.queue_wait_ms ?? 0; const duration = sample?.duration_ms ?? 0; + const level = analysis.levels.get(node.native_id) ?? 0; return { id: node.id, position: {x: position?.x ?? 0, y: position?.y ?? 0}, - data: {label:
{node.name}{node.id} + data: {label:
+
{taskflow_node_name(node.name)}层 {level + 1}
+ {node.name} {node.type} · W{sample?.worker_id ?? "--"} - 执行 {sample ? milliseconds(duration) : "未执行"} · 排队 {sample ? milliseconds(wait) : "--"}
}, + {Object.entries(node.attributes ?? {}).map(([key, value]) => + {key}:{value})} + 任务体 {sample ? milliseconds(duration) : "未执行"} · 排队 {sample ? milliseconds(wait) : "--"} + CPU {sample ? milliseconds(sample.cpu_duration_ms) : "--"} · 被抢占 {sample ? milliseconds(Math.max(0, duration - sample.cpu_duration_ms)) : "--"} + Observer {sample ? milliseconds(sample.observer_entry_ms + sample.observer_exit_ms) : "--"}
}, className: sample ? wait > duration && wait > .1 ? "taskflowNode taskflowNodeWaiting" : "taskflowNode taskflowNodeExecuted" : "taskflowNode" }; })); @@ -1225,10 +1338,33 @@ function Taskflow_Dag({graph, executions}: {graph: Taskflow_Graph_Trace; executi }); return () => { cancelled = true; }; }, [graph, executions]); - return
+ const copy_topology = async () => { + const native_ids = new Set(graph.nodes.map(node => node.native_id)); + try { + await navigator.clipboard.writeText(JSON.stringify({ + sequence: frame.sequence, + correlation_id: frame.correlation_id, + markers: frame.markers ?? {}, + graph, + executions: executions.filter(value => native_ids.has(value.native_id)) + }, null, 2)); + set_copy_state("已复制"); + window.setTimeout(() => set_copy_state("复制拓扑 JSON"), 1200); + } + catch { set_copy_state("复制失败"); } + }; + return
+
+ 未执行或条件跳过 + 已执行,执行时间为主 + 真实排队时间大于执行时间 +
+ +
-
; +
; } function Taskflow_Timeline({graph, executions}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]}) { @@ -1236,14 +1372,14 @@ function Taskflow_Timeline({graph, executions}: {graph: Taskflow_Graph_Trace; ex 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 end = Math.max(graph.finished_ms, ...rows.map(row => row.completed_ms), begin + .001); const duration = Math.max(.001, end - begin); - return
Worker 时间线横向位置按本帧真实 Observer 时间缩放;浅色段为估算就绪等待,亮色段为 on_entry → on_exit。
+ return
Worker 时间线浅色段为 Executor 等待,亮色段只表示真实任务体;Observer 开销在悬浮详情中单列。
{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
+ return
W{row.worker_id}{names.get(row.native_id) ?? row.node_id}
@@ -1289,6 +1425,21 @@ function Taskflow_Frame_Pane({plot}: {plot: Plot}) { }; const frame = response?.frames[frame_index]; const graph = frame?.graphs[graph_index]; + const graph_analysis = useMemo(() => graph && frame + ? taskflow_graph_analysis(graph, frame.executions) : null, [graph, frame]); + const paint_analysis = useMemo(() => { + const markers = frame?.markers; + if (!markers) return null; + const interval = (first: string, last: string) => + Math.max(0, (markers[last] ?? 0) - (markers[first] ?? 0)); + const total = interval("paint_started", "paint_finished"); + const frame_target = interval("paint_frame_target_started", "paint_frame_target_finished"); + const background = interval("paint_background_started", "paint_background_finished"); + const cache_targets = interval("paint_cache_targets_started", "paint_cache_targets_finished"); + const taskflow = interval("paint_taskflow_started", "paint_taskflow_finished"); + return {total, frame_target, background, cache_targets, taskflow, + coordination: Math.max(0, total - frame_target - background - cache_targets - taskflow)}; + }, [frame]); useEffect(() => set_graph_index(0), [frame?.sequence]); return
void load()}/>
@@ -1303,11 +1454,36 @@ function Taskflow_Frame_Pane({plot}: {plot: Plot}) { Executor {frame.worker_count} workers · 本帧 {frame.executions.length} 次任务执行
- {graph ? <>
+ {graph && graph_analysis ? <>
业务阶段
{graph.stage}
DAG 节点
{graph.nodes.length}
-
Topology 总耗时
{milliseconds(graph.finished_ms - graph.submitted_ms)}
+
Topology 墙钟
{milliseconds(graph_analysis.wall_time)}
+
任务体墙钟总和
{milliseconds(graph_analysis.execution_time)}
+
任务体实际 CPU
{milliseconds(graph_analysis.cpu_execution_time)}
+
Worker 被抢占/阻塞
{milliseconds(graph_analysis.descheduled_time)}
+
任务体墙钟并集
{milliseconds(graph_analysis.body_wall_time)}
+
Observer 独占墙钟
{milliseconds(graph_analysis.observer_wall_time)}
+
Observer entry 墙钟
{milliseconds(graph_analysis.observer_entry_time)}
+
Observer exit 墙钟
{milliseconds(graph_analysis.observer_exit_time)}
+
Observer 实际 CPU
{milliseconds(graph_analysis.observer_cpu_time)}
+
无任务墙钟
{milliseconds(graph_analysis.idle_wall_time)}
+
累计节点排队
{milliseconds(graph_analysis.queue_time)}
+
首次准入等待
{milliseconds(graph_analysis.initial_wait)}
+
内部调度空洞
{milliseconds(graph_analysis.internal_idle_time)}
+
Topology 完成尾部
{milliseconds(graph_analysis.completion_tail)}
+
入口队列峰值
{graph_analysis.peak_entry_queue}
+
实际最大并行
{graph_analysis.maximum_parallelism}
+
依赖层 / 并行层
{graph_analysis.layer_count} / {graph_analysis.parallel_layer_count}
+
最长执行节点
{graph_analysis.longest ? milliseconds(graph_analysis.longest.duration_ms) : "--"}
+ {graph.stage === "render_2d.paint" && paint_analysis ? <> +
同帧 Paint 总墙钟
{milliseconds(paint_analysis.total)}
+
帧目标准备
{milliseconds(paint_analysis.frame_target)}
+
背景填充
{milliseconds(paint_analysis.background)}
+
缓存目标准备
{milliseconds(paint_analysis.cache_targets)}
+
Taskflow marker 墙钟
{milliseconds(paint_analysis.taskflow)}
+
Paint 阶段衔接
{milliseconds(paint_analysis.coordination)}
+ : null}
状态
{graph.completed ? "完成" : "未完成"}
-
:
该帧没有 Taskflow 阶段
} +
:
该帧没有 Taskflow 阶段
} :
等待逐帧 Taskflow 样本输入 N 后捕获后续实际渲染帧;每帧独立绑定其 DAG 元信息和 Observer 结果。
}
; } diff --git a/webapp_gallery/src/styles.css b/webapp_gallery/src/styles.css index e731de4..fd5330b 100644 --- a/webapp_gallery/src/styles.css +++ b/webapp_gallery/src/styles.css @@ -194,13 +194,23 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .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; } +.taskflowDagSection { overflow: hidden; border: 1px solid #213653; border-radius: 10px; background: #07101c; } +.taskflowDagToolbar { display: flex; align-items: center; justify-content: space-between; gap: 12px; padding: 9px 11px; border-bottom: 1px solid #213653; background: #0c1727; } +.taskflowDagToolbar button { flex: none; padding: 6px 9px; color: #b9cce3; border: 1px solid #36516f; border-radius: 6px; background: #132238; cursor: pointer; } +.taskflowLegend { display: flex; align-items: center; flex-wrap: wrap; gap: 12px; color: #8298b4; font-size: 10px; } +.taskflowLegend span { display: inline-flex; align-items: center; gap: 5px; } +.taskflowLegend i { width: 10px; height: 10px; border: 1px solid #304766; border-radius: 3px; background: #0e1a2b; } +.taskflowLegend .taskflowLegendExecuted { border-color: #2c8e78; background: #0c201e; } +.taskflowLegend .taskflowLegendWaiting { border-color: #b9833e; background: #251b10; } +.taskflowDag { width: 100%; height: 620px; min-height: 420px; overflow: hidden; background: #07101c; } +.taskflowNode { width: 280px !important; min-height: 132px; padding: 11px !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 { display: grid; gap: 4px; min-width: 0; cursor: text; text-align: left; user-select: text !important; } +.taskflowNodeLabel header { display: flex; align-items: flex-start; justify-content: space-between; gap: 8px; } +.taskflowNodeLabel strong { color: #e1ecf9; font-size: 12px; line-height: 1.35; overflow-wrap: anywhere; } +.taskflowNodeLabel header i { flex: none; padding: 2px 5px; color: #74d8c0; border: 1px solid #2d675d; border-radius: 99px; font: 9px/1 ui-monospace, monospace; font-style: normal; } +.taskflowNodeLabel code { color: #6f89aa; font: 9px/1.25 ui-monospace, monospace; overflow-wrap: anywhere; white-space: normal; user-select: text !important; } .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; }