diff --git a/kernel/src/kernel/Frame_Scheduler.cpp b/kernel/src/kernel/Frame_Scheduler.cpp index f077e64..f0c96ae 100644 --- a/kernel/src/kernel/Frame_Scheduler.cpp +++ b/kernel/src/kernel/Frame_Scheduler.cpp @@ -1,16 +1,13 @@ #include "Frame_Scheduler.hpp" +#include #include #include #include #include -#include -#include #include #include -#include #include #include -#include #include namespace aethera { @@ -59,7 +56,8 @@ struct Frame_Scheduler::Private { periodic, once, cancel, - destroy + destroy, + stop }; struct Command { Command_Type type{}; @@ -71,10 +69,7 @@ struct Frame_Scheduler::Private { TimerWheel wheel{}; Scheduler_Clock::time_point origin{Scheduler_Clock::now()}; Scheduler_Clock::time_point last_advance{origin}; - std::mutex mutex{}; - std::condition_variable wake{}; - std::deque commands{}; - std::unordered_map> timers{}; + moodycamel::BlockingConcurrentQueue commands{256}; std::atomic_uint64_t next_id{1}; std::atomic_bool stopping{}; std::thread::id control_thread_id{}; @@ -85,26 +80,17 @@ struct Frame_Scheduler::Private { ~Private() { stopping.store(true, std::memory_order_release); thread.request_stop(); - wake.notify_all(); + enqueue({Command_Type::stop}); } void enqueue(Command command) { - { - std::lock_guard lock(mutex); - commands.push_back(std::move(command)); - } - wake.notify_one(); + if (!commands.enqueue(std::move(command))) std::terminate(); } std::shared_ptr make_timer(Handler handler) { if (!handler) throw std::invalid_argument("frame scheduler timer handler is empty"); const auto id = next_id.fetch_add(1, std::memory_order_relaxed); auto state = std::make_shared(this, id, std::move(handler)); - { - std::lock_guard lock(mutex); - /* 注册表只用于维持 TimerEvent 到异步 cancel/destroy 被处理为止的生命周期。 */ - timers.emplace(id, state); - } return state; } @@ -123,15 +109,10 @@ struct Frame_Scheduler::Private { wheel.schedule(&state.event, delay_ticks(state.next_deadline, now)); } - void process_commands(Scheduler_Clock::time_point now) { - std::deque local; - { - std::lock_guard lock(mutex); - local.swap(commands); - } - for (auto& command : local) { + void process_command(Command command, + Scheduler_Clock::time_point now) { auto& state = command.state; - if (!state) continue; + if (!state) return; switch (command.type) { case Command_Type::periodic: if (!state->alive.load(std::memory_order_acquire)) break; @@ -164,16 +145,21 @@ struct Frame_Scheduler::Private { case Command_Type::destroy: state->event.cancel(); state->mode = Timer::State::Mode::stopped; - { - std::lock_guard lock(mutex); - timers.erase(state->id); - } + break; + case Command_Type::stop: break; } if (command.completion) { try { command.completion->set_value(); } catch (...) {} } + } + + void process_commands(Scheduler_Clock::time_point now) { + Command command; + while (commands.try_dequeue(command)) { + process_command(std::move(command), now); + now = Scheduler_Clock::now(); } } @@ -232,21 +218,19 @@ struct Frame_Scheduler::Private { process_commands(now); advance_to(now); - std::unique_lock lock(mutex); if (stop.stop_requested() || stopping.load(std::memory_order_acquire)) break; - if (!commands.empty()) continue; - /* timers.size() 包含已取消但尚未 destroy 的句柄;用 wheel 查询决定实际等待。 */ constexpr ::Tick max_sleep_ticks = 10'000; /* 最多 1 s,便于可靠处理停止。 */ const ::Tick next = wheel.ticks_to_next_event(max_sleep_ticks); const auto deadline = last_advance + Frame_Scheduler::tick_duration * static_cast(next); - wake.wait_until(lock, deadline, [this, &stop] { - return !commands.empty() || stop.stop_requested() || - stopping.load(std::memory_order_acquire); - }); + Command command; + if (commands.wait_dequeue_timed( + command, std::max(deadline - Scheduler_Clock::now(), + Scheduler_Clock::duration::zero()))) + process_command(std::move(command), Scheduler_Clock::now()); } } catch (...) { diff --git a/kernel/src/kernel/Task_Graph.cpp b/kernel/src/kernel/Task_Graph.cpp index 8015613..eceebd4 100644 --- a/kernel/src/kernel/Task_Graph.cpp +++ b/kernel/src/kernel/Task_Graph.cpp @@ -60,6 +60,13 @@ struct Task_Graph::Private { if (source.child) source.child->collect_nodes(module_id, result); } } + + void bind_frame(Render_Frame* frame) noexcept { + for (auto& node : nodes) { + node.task.data(frame); + if (node.child) node.child->bind_frame(frame); + } + } }; struct Task_Node::Private { @@ -79,6 +86,11 @@ std::vector Task_Graph_Access::nodes( graph.d->collect_nodes({}, result); return result; } + +void Task_Graph_Access::bind_frame(Task_Graph& graph, + Render_Frame* frame) noexcept { + graph.d->bind_frame(frame); +} } Task_Node::Task_Node() = default; diff --git a/kernel/src/kernel/Task_Graph_Internal.hpp b/kernel/src/kernel/Task_Graph_Internal.hpp index 966e51c..5035554 100644 --- a/kernel/src/kernel/Task_Graph_Internal.hpp +++ b/kernel/src/kernel/Task_Graph_Internal.hpp @@ -12,6 +12,7 @@ struct Task_Graph_Access { [[nodiscard]] static void* native_storage(Task_Graph& graph) noexcept; [[nodiscard]] static std::vector nodes( const Task_Graph& graph); + static void bind_frame(Task_Graph& graph, Render_Frame* frame) noexcept; }; /* Task_Graph.cpp 到全局 Kernel Executor 的内部桥;不得在模块公共头文件中声明。 */ void corun_taskflow_until(std::function predicate); diff --git a/kernel/src/kernel/Taskflow_Frame_Access.hpp b/kernel/src/kernel/Taskflow_Frame_Access.hpp index c028fb9..19639a3 100644 --- a/kernel/src/kernel/Taskflow_Frame_Access.hpp +++ b/kernel/src/kernel/Taskflow_Frame_Access.hpp @@ -9,15 +9,16 @@ struct Taskflow_Graph_Token { Render_Frame* frame{}; /* DAG 运行记录所属外部帧;帧完成前有效。 */ std::size_t index{}; /* 帧内 graph 记录索引。 */ }; - +struct Taskflow_Graph_Timing { + std::chrono::steady_clock::time_point executor_finished{}; + std::chrono::steady_clock::time_point runtime_bookkeeping_finished{}; + std::chrono::steady_clock::time_point runtime_publish_finished{}; + std::chrono::steady_clock::time_point observer_unbind_finished{}; +}; struct Taskflow_Frame_Access { using Clock = std::chrono::steady_clock; static void begin_capture(Render_Frame& frame, std::size_t workers); static void finish_capture(Render_Frame& frame) noexcept; - [[nodiscard]] static bool acquire_writer(Render_Frame& frame) noexcept; - static void release_writer(Render_Frame& frame) noexcept; - [[nodiscard]] static bool contains_task(const Render_Frame& frame, - std::uint64_t native_id) noexcept; [[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, @@ -32,6 +33,7 @@ struct Taskflow_Frame_Access { 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; + static void finish_graph(Taskflow_Graph_Token token, + const Taskflow_Graph_Timing& timing) noexcept; }; } diff --git a/kernel/src/kernel/double_buffer/Dependency_Graph_Storage.hpp b/kernel/src/kernel/double_buffer/Dependency_Graph_Storage.hpp index c5fd6dd..4736b39 100644 --- a/kernel/src/kernel/double_buffer/Dependency_Graph_Storage.hpp +++ b/kernel/src/kernel/double_buffer/Dependency_Graph_Storage.hpp @@ -369,6 +369,7 @@ struct Root::Builder { auto validate_result = validate(); if (!validate_result) return std::unexpected(validate_result.error()); auto private_data = std::make_unique(object->memory_resource()); + *private_data->pending = prop; *private_data->current = prop; buffer_storage.commit(private_data->buffer_storage); dependency_graph_storage.commit(private_data->dependency_graph_storage); @@ -386,5 +387,11 @@ protected: [[nodiscard]] static Private_Access private_access(Target* target) noexcept { return Private_Access{static_cast(*target->d)}; } + /* Builder 内部分派读取目标已提交 Prop;不经过 Impl 对外的 pending 读取接口。 */ + template + [[nodiscard]] static const auto& current_prop(const Target* target) noexcept { + const auto& private_data = static_cast(*target->d); + return private_data.template current_prop(); + } }; } diff --git a/kernel/src/kernel/double_buffer/mechanism.hpp b/kernel/src/kernel/double_buffer/mechanism.hpp index de3a683..1bef015 100644 --- a/kernel/src/kernel/double_buffer/mechanism.hpp +++ b/kernel/src/kernel/double_buffer/mechanism.hpp @@ -470,13 +470,16 @@ template using Prop_Value = typename Prop_By_Tag::Type; template concept State_Callback_For = State_Tag_In && State_Chain_Matches && std::invocable&>; +template +struct State_Callback_Entry { + std::mutex publish_lock{}; /* 仅串行化该状态层的回调安装、清除与发布。 */ + std::function callback{}; /* 当前状态层唯一的发布回调。 */ +}; template struct State_Callback_Tuple; template struct State_Callback_Tuple> { - using Type = std::tuple - ... - >; + using Type = std::tuple...>; }; template struct State_Callback_Storage { @@ -493,24 +496,27 @@ public: void set(Callback&& callback) { constexpr auto value_index = index(); using Layer = std::tuple_element_t; - std::get < value_index > (callbacks) = std::function < void(const Layer &) > - { - std::forward(callback) - }; + auto& entry = std::get(callbacks); + std::lock_guard guard(entry.publish_lock); + entry.callback = std::function{std::forward(callback)}; } template Tag> void clear() { - std::get < index() > (callbacks) = {}; + auto& entry = std::get()>(callbacks); + std::lock_guard guard(entry.publish_lock); + entry.callback = {}; } - template Tag> - void notify(const State_T& state) { + template Tag, typename Publish> + requires std::same_as, State_Value> + void publish(Publish&& publish) { constexpr auto value_index = index(); - using Layer = std::tuple_element_t; - auto& callback = std::get < value_index > (callbacks); - if (callback) callback(static_cast(state)); + auto& entry = std::get(callbacks); + std::lock_guard guard(entry.publish_lock); + auto state = std::invoke(std::forward(publish)); + if (entry.callback) entry.callback(state); } }; -// 每个 State Tag 独立保存回调;字段更新不会隐式通知,发布频率由拥有该状态的内部机制显式调用 notify 控制。 +// 每个 State Tag 独立保存回调与发布锁;字段更新不会隐式通知,发布频率由拥有该状态的内部机制显式控制。 template struct Buffer_Value_By_Tag; template @@ -706,6 +712,16 @@ struct Root { /* CRTP 默认:Builder 挂接最终 Private 后按基类到派生类绑定最终对象类型;Root 层不处理。 */ template void bind_private_crtp(Object*) {} + protected: + /* 内部实现按目标类型读取已提交 Prop;不经过 Impl 对外的 pending 读取接口。 */ + template + requires requires(const typename Target::Private& private_data) { + private_data.template current_prop(); + } + [[nodiscard]] static const auto& read_current_prop(const Target* target) noexcept { + const auto& private_data = static_cast(*target->d); + return private_data.template current_prop(); + } }; using Buffers = std::tuple<>; using Dependency_Graph_Types = std::tuple<>; diff --git a/kernel/src/kernel/double_buffer/model.hpp b/kernel/src/kernel/double_buffer/model.hpp index 5eba354..0a71a9b 100644 --- a/kernel/src/kernel/double_buffer/model.hpp +++ b/kernel/src/kernel/double_buffer/model.hpp @@ -133,17 +133,32 @@ struct Impl : Obj { using Base::Base; }; using Attached_Object = Obj; - struct Private : Obj::Private, Publish_Double_Buffer { + struct Private : Obj::Private, Commit_Double_Buffer { using Allocator = std::pmr::polymorphic_allocator; + mutable Lock advance_lock; /* 仅保护全部双缓冲在 advance/publish 边界的交换。 */ Commit_Double_Buffer state; detail::State_Callback_Storage state_callbacks; detail::Buffer_Storage buffer_storage; detail::Dependency_Graph_Storage dependency_graph_storage; detail::Mpmc_Triple_Buffer_Storage_Set mpmc_triple_buffer_storage; - explicit Private(std::pmr::memory_resource* resource) : Publish_Double_Buffer(std::allocator_arg, Allocator{resource}), + explicit Private(std::pmr::memory_resource* resource) : Commit_Double_Buffer(std::allocator_arg, Allocator{resource}), state(std::allocator_arg, Allocator{resource}), buffer_storage(std::allocator_arg, Allocator{resource}), dependency_graph_storage(std::allocator_arg, Allocator{resource}) {} + template requires detail::Prop_Tag_In> + [[nodiscard]] const auto& current_prop() const noexcept { + using Layer = detail::Prop_Value; + return static_cast(*this->current); + } + template Tag> + void publish_state() { + using Layer = detail::State_Value; + state_callbacks.template publish([&] { + std::lock_guard guard(advance_lock); + state.advance(); + return Layer{static_cast(*state.current)}; + }); + } }; private: friend struct Root::Builder; @@ -152,7 +167,6 @@ private: this->template bind_object_crtp(); } friend struct Root; - mutable Lock lock; [[nodiscard]] Private& data() noexcept { return static_cast(*this->d); } @@ -169,13 +183,13 @@ private: } template void before_prop_set(Member Owner::* member) { - Prop_Access props{*data().current}; + Prop_Access props{*data().pending}; auto callback = [&](auto& private_data) { private_data.before_prop_set(this, member, props); }; walk_private(callback); } template void after_prop_set(Member Owner::* member) { - Prop_Access props{*data().current}; + Prop_Access props{*data().pending}; auto callback = [&](auto& private_data) { private_data.after_prop_set(this, member, props); }; walk_private(callback); } @@ -236,16 +250,12 @@ private: }; walk_private(callback); } - void advance_unlocked() { - before_advance(); - // State 从 pending 向内部 current 提交;Prop 从内部 current 向 pending 发布,两者推进后的写入侧都同步为最新基线。 - static_cast&>(data()).advance(); + void exchange() { + // Prop 与 State 均从外部 pending 提交到内部 current,交换后 pending 同步为最新编辑基线。 + static_cast&>(data()).advance(); data().state.advance(); data().buffer_storage.advance(); data().dependency_graph_storage.advance(); - after_advance(); - // Commit_Double_Buffer::advance() 已为 pending 建立 current 基线;after_advance - // 只向 pending 写本轮派生状态,禁止回写正在供外部查询的 current。 } public: template Tag> @@ -310,7 +320,6 @@ public: // 依赖图只允许在编辑侧修改;回调完成后先验证 DAG,再把新结构提交到 pending。 template ... Tags, detail::Dependency_Graph_Edit_Callback_For Callback> std::expected edit_dependency_graph(Callback&& callback) { - std::lock_guard guard(lock); return data().dependency_graph_storage.template edit(std::forward(callback)); } std::expected validate_dependency_graph() const { @@ -326,47 +335,41 @@ public: } template requires detail::Prop_Member && requires(Prop& prop, Value&& value) { prop.*Member = std::forward(value); } Impl& set(Value&& value) { - std::lock_guard guard(lock); before_prop_set(Member); - data().current->*Member = std::forward(value); + data().pending->*Member = std::forward(value); after_prop_set(Member); emit_prop_dependencies(); return *this; } template requires detail::Prop_Dependency_Source [[nodiscard]] auto get() const { - std::lock_guard guard(lock); - return data().current->*Member; + return data().pending->*Member; } template requires detail::Prop_Tag_In> [[nodiscard]] const auto& read_prop() const noexcept { using Layer = detail::Prop_Value; - return static_cast(*data().current); + return static_cast(*data().pending); } template requires (sizeof...(Members) > 0) && (detail::Prop_Member && ...) && std::invocable> void update_prop(Callback&& callback) { - std::lock_guard guard(lock); (before_prop_set(Members), ...); - std::invoke(std::forward(callback), Prop_Access{*data().current}); + std::invoke(std::forward(callback), Prop_Access{*data().pending}); (after_prop_set(Members), ...); (emit_prop_dependencies(), ...); } template Tag, detail::State_Callback_For Callback> void set_state_callback(Callback&& callback) { - std::lock_guard guard(lock); data().state_callbacks.template set(std::forward(callback)); } template Tag> void clear_state_callback() { - std::lock_guard guard(lock); data().state_callbacks.template clear(); } template Tag, detail::State_Callback_For Callback> void access_state(Callback&& callback) const { - std::lock_guard guard(lock); using Layer = detail::State_Value; std::invoke(std::forward(callback), static_cast(*data().state.current)); } @@ -376,14 +379,9 @@ public: using Layer = detail::State_Value; return static_cast(*data().state.current); } - template Tag> - void notify_state() { - data().state_callbacks.template notify(*data().state.current); - } // State 写入 pending,随后同时发出成员级与状态层级依赖信号;真正进入 current 发生在 advance/process 边界。 template Value> void update_state(Value&& value) { - std::lock_guard guard(lock); before_state_set(Member); data().state.pending->*Member = std::forward(value); after_state_set(Member); @@ -395,7 +393,6 @@ public: (detail::State_Member && ...) && std::invocable> void update_state(Callback&& callback) { - std::lock_guard guard(lock); (before_state_set(Members), ...); std::invoke(std::forward(callback), State_Access{*data().state.pending}); (after_state_set(Members), ...); @@ -403,38 +400,38 @@ public: } template Callback> void process(Callback&& callback) { - std::lock_guard guard(lock); - advance_unlocked(); + advance(); data().process(this, std::forward(callback)); } void advance() { - std::lock_guard guard(lock); - advance_unlocked(); + before_advance(); + { + std::lock_guard guard(data().advance_lock); + exchange(); + } + after_advance(); + // Commit_Double_Buffer::advance() 已为 pending 建立 current 基线;after_advance + // 只向 pending 写本轮派生状态,禁止回写正在供外部查询的 current。 } template Tag> void publish_state() { - std::lock_guard guard(lock); - data().state.advance(); - data().state_callbacks.template notify(*data().state.current); + data().template publish_state(); } template Tag, auto... Members, typename Callback> requires (sizeof...(Members) > 0) && (detail::State_Member && ...) && std::invocable> void publish_state(Callback&& callback) { - std::lock_guard guard(lock); (before_state_set(Members), ...); std::invoke(std::forward(callback), State_Access{*data().state.pending}); (after_state_set(Members), ...); (emit_state_dependencies(), ...); - data().state.advance(); - data().state_callbacks.template notify(*data().state.current); + publish_state(); } template Callback> void advance(Callback&& callback) { - std::lock_guard guard(lock); - advance_unlocked(); + advance(); std::invoke(std::forward(callback), std::as_const(*data().state.current)); } }; diff --git a/kernel/src/kernel/frame.cpp b/kernel/src/kernel/frame.cpp index cef80b2..2334633 100644 --- a/kernel/src/kernel/frame.cpp +++ b/kernel/src/kernel/frame.cpp @@ -8,10 +8,8 @@ #include #include #include -#include #include #include -#include #include #include namespace aethera { @@ -31,11 +29,8 @@ struct Render_Frame::Private { std::array markers{}; /* 每种时间点首次出现时的 elapsed_ns 加一编码。 */ std::array measurements{}; /* 每种原始耗时首次记录值的加一编码。 */ std::atomic_bool taskflow_trace_requested{}; /* 本次逻辑帧是否请求原生 Taskflow 捕获。 */ - std::atomic_bool taskflow_trace_capturing{}; /* 原生 Observer 是否仍可写入本帧。 */ - std::atomic_size_t taskflow_trace_writers{}; /* 正在完成 on_exit 写入的 worker 数。 */ std::vector taskflow_workers{}; /* 按 Executor worker 隔离的单写者时间线。 */ - std::mutex taskflow_graph_mutex{}; /* 只在按需捕获时保护跨阶段 DAG 追加与完成标记。 */ - std::vector taskflow_graphs{}; /* 本帧主动 run 的业务 DAG 元信息。 */ + std::vector taskflow_graphs{}; /* 内部流水线追加完成后,由外部一次性消费。 */ }; Render_Frame::Render_Frame(Frame_Identity identity) : d(std::make_unique()) { begin(identity); @@ -50,13 +45,8 @@ void Render_Frame::begin(Frame_Identity identity) noexcept { for (auto& measurement : d->measurements) measurement.store(0, std::memory_order_relaxed); d->taskflow_trace_requested.store(false, std::memory_order_relaxed); - d->taskflow_trace_capturing.store(false, std::memory_order_relaxed); - d->taskflow_trace_writers.store(0, std::memory_order_relaxed); d->taskflow_workers.clear(); - { - std::lock_guard lock(d->taskflow_graph_mutex); - d->taskflow_graphs.clear(); - } + d->taskflow_graphs.clear(); d->markers[static_cast(Frame_Trace_Marker::created)].store(encode_present_value(0), std::memory_order_relaxed); } Render_Frame::~Render_Frame() = default; @@ -273,7 +263,7 @@ bool Render_Frame::taskflow_trace_requested() const noexcept { return d->taskflow_trace_requested.load(std::memory_order_acquire); } -Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { +Taskflow_Frame_Trace Render_Frame::take_taskflow_trace() { Taskflow_Frame_Trace result{}; result.identity = d->identity; result.created_time_unix_ns = d->created_time_unix_ns; @@ -284,12 +274,12 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { 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; + result.graphs = std::move(d->taskflow_graphs); + for (auto& worker : d->taskflow_workers) { + for (auto& task : worker.tasks) + result.tasks.push_back(std::move(task)); + worker.tasks.clear(); } - for (const auto& worker : d->taskflow_workers) - result.tasks.insert(result.tasks.end(), worker.tasks.begin(), worker.tasks.end()); std::ranges::sort(result.tasks, {}, &Taskflow_Task_Trace::started_ms); std::erase_if(result.tasks, [&](const Taskflow_Task_Trace& task) { @@ -375,6 +365,74 @@ Taskflow_Frame_Trace Render_Frame::taskflow_trace() const { task.queue_wait_ms = std::max(0.0, task.entered_ms - ready); previous_execution_finished[task.native_id] = task.completed_ms; } + std::unordered_set used_native_ids; + for (const auto& graph : result.graphs) + for (const auto& node : graph.nodes) used_native_ids.insert(node.native_id); + std::uint64_t diagnostic_native_id = std::numeric_limits::max(); + const auto allocate_diagnostic_id = [&]() { + while (used_native_ids.contains(diagnostic_native_id)) --diagnostic_native_id; + const auto value = diagnostic_native_id--; + used_native_ids.insert(value); + return value; + }; + for (auto& graph : result.graphs) { + if (!graph.completed || graph.finished_ms <= 0.0) continue; + std::unordered_set graph_native_ids; + for (const auto& node : graph.nodes) graph_native_ids.insert(node.native_id); + double last_completed = graph.submitted_ms; + for (const auto& task : result.tasks) + if (graph_native_ids.contains(task.native_id)) + last_completed = std::max(last_completed, task.completed_ms); + std::vector predecessors; + for (const auto& node : graph.nodes) + if (node.successors.empty()) predecessors.push_back(node.native_id); + struct Diagnostic_Phase { + const char* name; + const char* label; + double finished_ms; + }; + const Diagnostic_Phase phases[] = { + {"taskflow.executor.finalize", "Executor topology 收尾", graph.executor_finished_ms}, + {"taskflow.runtime.bookkeeping", "运行计数更新", graph.runtime_bookkeeping_finished_ms}, + {"taskflow.runtime.publish_state", "运行时状态发布", graph.runtime_publish_finished_ms}, + {"taskflow.observer.unbind_frame", "Observer 帧绑定清除", graph.observer_unbind_finished_ms}, + {"taskflow.graph.finish", "Graph 完成记录", graph.finished_ms} + }; + double started_ms = last_completed; + for (const auto& phase : phases) { + if (phase.finished_ms <= 0.0) continue; + const auto finished_ms = std::max(started_ms, phase.finished_ms); + const auto native_id = allocate_diagnostic_id(); + Taskflow_Graph_Trace::Node node{}; + node.native_id = native_id; + node.node_id = (graph.taskflow_name.empty() ? graph.stage : graph.taskflow_name) + + "/@diagnostic/" + phase.name; + node.name = phase.name; + node.type = "diagnostic"; + node.predecessors = predecessors; + node.attributes.emplace_back("diagnostic", "runtime_tail"); + node.attributes.emplace_back("label", phase.label); + for (const auto predecessor : predecessors) { + const auto found = std::ranges::find( + graph.nodes, predecessor, &Taskflow_Graph_Trace::Node::native_id); + if (found != graph.nodes.end()) found->successors.push_back(native_id); + } + graph.nodes.push_back(std::move(node)); + Taskflow_Task_Trace task{}; + task.native_id = native_id; + task.worker_id = result.worker_count; + task.ready_ms = started_ms; + task.entered_ms = started_ms; + task.started_ms = started_ms; + task.finished_ms = finished_ms; + task.completed_ms = finished_ms; + task.duration_ms = std::max(0.0, finished_ms - started_ms); + result.tasks.push_back(std::move(task)); + predecessors.assign(1, native_id); + started_ms = finished_ms; + } + } + std::ranges::sort(result.tasks, {}, &Taskflow_Task_Trace::started_ms); return result; } @@ -383,50 +441,12 @@ void detail::Taskflow_Frame_Access::begin_capture(Render_Frame& frame, std::size data.taskflow_workers.clear(); data.taskflow_workers.resize(workers); for (auto& worker : data.taskflow_workers) worker.tasks.reserve(64); - { - std::lock_guard lock(data.taskflow_graph_mutex); - data.taskflow_graphs.clear(); - } - data.taskflow_trace_writers.store(0, std::memory_order_relaxed); - data.taskflow_trace_capturing.store(true, std::memory_order_release); + data.taskflow_graphs.clear(); } void detail::Taskflow_Frame_Access::finish_capture(Render_Frame& frame) noexcept { - auto& data = *frame.d; - data.taskflow_trace_capturing.store(false, std::memory_order_release); - while (data.taskflow_trace_writers.load(std::memory_order_acquire) != 0) - std::this_thread::yield(); -} - -bool detail::Taskflow_Frame_Access::acquire_writer(Render_Frame& frame) noexcept { - auto& data = *frame.d; - if (!data.taskflow_trace_capturing.load(std::memory_order_acquire)) return false; - data.taskflow_trace_writers.fetch_add(1, std::memory_order_acq_rel); - if (data.taskflow_trace_capturing.load(std::memory_order_acquire)) return true; - data.taskflow_trace_writers.fetch_sub(1, std::memory_order_release); - return false; -} - -void detail::Taskflow_Frame_Access::release_writer(Render_Frame& frame) noexcept { - frame.d->taskflow_trace_writers.fetch_sub(1, std::memory_order_release); -} - -bool detail::Taskflow_Frame_Access::contains_task( - const Render_Frame& frame, std::uint64_t native_id) noexcept { - try { - std::lock_guard lock(frame.d->taskflow_graph_mutex); - return std::ranges::any_of( - frame.d->taskflow_graphs, - [&](const Taskflow_Graph_Trace& graph) { - return std::ranges::any_of( - graph.nodes, [&](const Taskflow_Graph_Trace::Node& node) { - return node.native_id == native_id; - }); - }); - } - catch (...) { - return false; - } + (void)frame; + /* run completion 已保证最后一个 Observer on_exit 返回;这里只标记消费边界,不同步等待。 */ } std::size_t detail::Taskflow_Frame_Access::append_task( @@ -490,20 +510,27 @@ detail::Taskflow_Graph_Token detail::Taskflow_Frame_Access::begin_graph( graph.nodes = detail::Task_Graph_Access::nodes(taskflow); graph.submitted_ms = std::chrono::duration( Clock::now() - frame.d->created_at).count(); - std::lock_guard lock(frame.d->taskflow_graph_mutex); frame.d->taskflow_graphs.push_back(std::move(graph)); return {&frame, frame.d->taskflow_graphs.size() - 1}; } void detail::Taskflow_Frame_Access::finish_graph( - detail::Taskflow_Graph_Token token) noexcept { + detail::Taskflow_Graph_Token token, + const detail::Taskflow_Graph_Timing& timing) noexcept { if (!token.frame) return; auto& data = *token.frame->d; - std::lock_guard lock(data.taskflow_graph_mutex); + const auto elapsed_ms = [&](Clock::time_point value) { + return value == Clock::time_point{} ? 0.0 + : std::chrono::duration(value - data.created_at).count(); + }; if (token.index >= data.taskflow_graphs.size()) return; + const auto finished = Clock::now(); auto& graph = data.taskflow_graphs[token.index]; - graph.finished_ms = std::chrono::duration( - Clock::now() - data.created_at).count(); + graph.executor_finished_ms = elapsed_ms(timing.executor_finished); + graph.runtime_bookkeeping_finished_ms = elapsed_ms(timing.runtime_bookkeeping_finished); + graph.runtime_publish_finished_ms = elapsed_ms(timing.runtime_publish_finished); + graph.observer_unbind_finished_ms = elapsed_ms(timing.observer_unbind_finished); + graph.finished_ms = elapsed_ms(finished); graph.completed = true; } } diff --git a/kernel/src/kernel/frame.hpp b/kernel/src/kernel/frame.hpp index 10fa914..496ad21 100644 --- a/kernel/src/kernel/frame.hpp +++ b/kernel/src/kernel/frame.hpp @@ -87,6 +87,10 @@ struct Taskflow_Graph_Trace { }; std::vector nodes{}; /* DAG 构造时登记的节点与依赖元信息。 */ double submitted_ms{}; /* 相对帧创建时刻的 run 提交时间。 */ + double executor_finished_ms{}; /* Executor 同步返回或异步 completion 进入的时间。 */ + double runtime_bookkeeping_finished_ms{}; /* Taskflow 运行计数更新完成时间。 */ + double runtime_publish_finished_ms{}; /* Taskflow 运行时状态发布完成时间。 */ + double observer_unbind_finished_ms{}; /* Node 帧绑定清除完成时间。 */ double finished_ms{}; /* 同步返回或异步 topology 完成时间。 */ bool completed{}; /* 对应 topology 是否已经结束。 */ }; @@ -130,7 +134,7 @@ public: [[nodiscard]] Frame_Statistics_Sample statistics(Frame_Dimension dimension) const; void request_taskflow_trace() noexcept; [[nodiscard]] bool taskflow_trace_requested() const noexcept; - [[nodiscard]] Taskflow_Frame_Trace taskflow_trace() const; + [[nodiscard]] Taskflow_Frame_Trace take_taskflow_trace(); protected: /* 物理帧槽再次承载新逻辑帧时,重建其唯一身份和诊断时间原点。 */ void begin(Frame_Identity identity) noexcept; diff --git a/kernel/src/kernel/render_common.cpp b/kernel/src/kernel/render_common.cpp index 4516f80..fcc7efc 100644 --- a/kernel/src/kernel/render_common.cpp +++ b/kernel/src/kernel/render_common.cpp @@ -10,7 +10,6 @@ #include #include #include -#include #include #include #include @@ -37,6 +36,21 @@ tf::Taskflow& native_taskflow(Task_Graph& graph) noexcept { return *static_cast( detail::Task_Graph_Access::native_storage(graph)); } +Render_Frame* task_frame(const tf::TaskView& task) noexcept { + /* + * Taskflow v4.1.0 ABI: TaskView 的 offset 0 是唯一的 const Node&,Task 的 + * offset 0 是唯一的 Node*。这里把同一 Node 转成公开 Task 句柄,只读取 + * Task::data(),避免修改第三方代码。版本升级导致句柄尺寸变化时立即编译失败。 + */ + static_assert(sizeof(tf::TaskView) == sizeof(tf::Node*)); + static_assert(alignof(tf::TaskView) == alignof(tf::Node*)); + static_assert(sizeof(tf::Task) == sizeof(tf::Node*)); + static_assert(alignof(tf::Task) == alignof(tf::Node*)); + auto* node = *reinterpret_cast(std::addressof(task)); + tf::Task handle{}; + *reinterpret_cast(std::addressof(handle)) = node; + return static_cast(handle.data()); +} #if defined(_WIN32) std::uint64_t thread_cpu_ns(HANDLE thread) noexcept { FILETIME created{}, exited{}, kernel{}, user{}; @@ -156,8 +170,6 @@ private: std::vector worker_cpu_starts; std::unique_ptr worker_statistics; std::size_t worker_statistics_count{}; - std::unordered_map trace_tasks{}; - std::shared_mutex trace_mutex{}; /* 按原生节点身份把并发 Taskflow 归属到各自 Frame。 */ std::uint64_t worker_occupation_limit_ns{}; /* 单节点连续非 CPU 等待 Worker 的上限。 */ Task_Overrun_Action worker_overrun_action{}; /* 节点超过占用上限后的处置策略。 */ std::atomic_bool watchdog_stopping{}; /* Watchdog 生命周期停止标志。 */ @@ -433,21 +445,7 @@ public: record.native_id = static_cast(task.hash_value()); record.type = task.type(); worker_starts.push_back(std::move(record)); - Render_Frame* frame{}; - /* - * 每次被追踪的 Taskflow run 在提交前登记 native node -> Frame。 - * 因而 Render(N+1) 与外接 H264(N) 可以同时被 Observer 正确归属, - * 不再使用“全局唯一当前 Frame”的假设。帧租约从 on_entry 持续到 - * on_exit,避免 graph 退役后迟到的 Observer 写入访问已复用 Frame。 - */ - { - std::shared_lock trace_guard(trace_mutex); - const auto found = trace_tasks.find( - static_cast(task.hash_value())); - if (found != trace_tasks.end() && - detail::Taskflow_Frame_Access::acquire_writer(*found->second)) - frame = found->second; - } + Render_Frame* frame = task_frame(task); worker_starts.back().frame = frame; const auto active_depth = worker_state.active_depth.fetch_add( 1, std::memory_order_relaxed) + 1; @@ -568,7 +566,6 @@ public: 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( @@ -647,40 +644,7 @@ public: detail::Taskflow_Frame_Access::begin_capture(frame, workers); return true; } - void activate_graph(Render_Frame& frame, Task_Graph& graph) { - if (!frame.taskflow_trace_requested()) return; - const auto nodes = detail::Task_Graph_Access::nodes(graph); - std::unique_lock guard(trace_mutex); - for (const auto& node : nodes) { - const auto [found, inserted] = trace_tasks.emplace( - node.native_id, &frame); - if (!inserted && found->second != &frame) - throw std::logic_error( - "concurrent traced runs share the same Taskflow node"); - } - } - void deactivate_graph(Render_Frame& frame, Task_Graph& graph) noexcept { - if (!frame.taskflow_trace_requested()) return; - try { - const auto nodes = detail::Task_Graph_Access::nodes(graph); - std::unique_lock guard(trace_mutex); - for (const auto& node : nodes) { - const auto found = trace_tasks.find(node.native_id); - if (found != trace_tasks.end() && found->second == &frame) - trace_tasks.erase(found); - } - } - catch (...) {} - } void finish_trace(Render_Frame& frame) noexcept { - try { - std::unique_lock guard(trace_mutex); - for (auto it = trace_tasks.begin(); it != trace_tasks.end();) { - if (it->second == &frame) it = trace_tasks.erase(it); - else ++it; - } - } - catch (...) {} detail::Taskflow_Frame_Access::finish_capture(frame); } void write_state(Task_Runtime_State& state, std::size_t workers, std::size_t active_topologies) const { @@ -873,13 +837,84 @@ private: } void publish_state() { if (!state_callback_enabled.load(std::memory_order_acquire)) return; - std::lock_guard guard(state_mutex); - observer->write_state(state, executor->num_workers(), executor->num_topologies()); - state.active_taskflow_count = active_taskflows.load(std::memory_order_relaxed); - state.peak_active_taskflow_count = peak_active_taskflows.load(std::memory_order_relaxed); - state.completed_taskflow_count = completed_taskflows.load(std::memory_order_relaxed); - state.failed_taskflow_count = failed_taskflows.load(std::memory_order_relaxed); - state_callbacks.template notify(state); + Task_Runtime_State published; + { + std::lock_guard guard(state_mutex); + observer->write_state(state, executor->num_workers(), executor->num_topologies()); + state.active_taskflow_count = active_taskflows.load(std::memory_order_relaxed); + state.peak_active_taskflow_count = peak_active_taskflows.load(std::memory_order_relaxed); + state.completed_taskflow_count = completed_taskflows.load(std::memory_order_relaxed); + state.failed_taskflow_count = failed_taskflows.load(std::memory_order_relaxed); + published = state; + } + state_callbacks.template publish( + [published = std::move(published)] { return published; }); + } + std::uint64_t run_synchronous( + Task_Graph& taskflow, detail::Taskflow_Graph_Timing* timing) { + ensure_executor(); + auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1; + auto peak = peak_active_taskflows.load(std::memory_order_relaxed); + while (peak < active && !peak_active_taskflows.compare_exchange_weak(peak, active, std::memory_order_relaxed)) {} + const auto start = std::chrono::steady_clock::now(); + try { + if (const auto* worker = executor->this_worker()) { + observer->pause_worker(worker->id()); + try { + executor->corun(native_taskflow(taskflow)); + } + catch (...) { + observer->resume_worker(worker->id()); + throw; + } + observer->resume_worker(worker->id()); + } + else executor->run(native_taskflow(taskflow)).get(); + } + catch (...) { + const auto executor_finished = std::chrono::steady_clock::now(); + active_taskflows.fetch_sub(1, std::memory_order_relaxed); + failed_taskflows.fetch_add(1, std::memory_order_relaxed); + if (timing) timing->runtime_bookkeeping_finished = std::chrono::steady_clock::now(); + publish_state(); + if (timing) { + timing->executor_finished = executor_finished; + timing->runtime_publish_finished = std::chrono::steady_clock::now(); + } + throw; + } + const auto executor_finished = std::chrono::steady_clock::now(); + const auto elapsed = static_cast( + std::chrono::duration_cast( + executor_finished - start).count()); + active_taskflows.fetch_sub(1, std::memory_order_relaxed); + completed_taskflows.fetch_add(1, std::memory_order_relaxed); + if (timing) timing->runtime_bookkeeping_finished = std::chrono::steady_clock::now(); + publish_state(); + if (timing) { + timing->executor_finished = executor_finished; + timing->runtime_publish_finished = std::chrono::steady_clock::now(); + } + return elapsed; + } + void run_asynchronous( + Task_Graph& taskflow, bool capture_timing, + std::function completion) { + ensure_executor(); + auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1; + auto peak = peak_active_taskflows.load(std::memory_order_relaxed); + while (peak < active && !peak_active_taskflows.compare_exchange_weak( + peak, active, std::memory_order_relaxed)) {} + executor->run(native_taskflow(taskflow), [this, capture_timing, completion = std::move(completion)]() mutable { + detail::Taskflow_Graph_Timing timing{}; + if (capture_timing) timing.executor_finished = std::chrono::steady_clock::now(); + active_taskflows.fetch_sub(1, std::memory_order_relaxed); + completed_taskflows.fetch_add(1, std::memory_order_relaxed); + if (capture_timing) timing.runtime_bookkeeping_finished = std::chrono::steady_clock::now(); + publish_state(); + if (capture_timing) timing.runtime_publish_finished = std::chrono::steady_clock::now(); + completion(std::move(timing)); + }); } public: Task_Resource() : memory(std::pmr::get_default_resource()) {} @@ -905,48 +940,14 @@ public: failed_taskflows.store(0, std::memory_order_relaxed); } std::uint64_t run(Task_Graph& taskflow) { - ensure_executor(); - auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1; - auto peak = peak_active_taskflows.load(std::memory_order_relaxed); - while (peak < active && !peak_active_taskflows.compare_exchange_weak(peak, active, std::memory_order_relaxed)) {} - auto start = std::chrono::steady_clock::now(); - try { - if (const auto* worker = executor->this_worker()) { - observer->pause_worker(worker->id()); - try { - executor->corun(native_taskflow(taskflow)); - } - catch (...) { - observer->resume_worker(worker->id()); - throw; - } - observer->resume_worker(worker->id()); - } - else executor->run(native_taskflow(taskflow)).get(); - } - catch (...) { - active_taskflows.fetch_sub(1, std::memory_order_relaxed); - failed_taskflows.fetch_add(1, std::memory_order_relaxed); - publish_state(); - throw; - } - auto elapsed = static_cast(std::chrono::duration_cast(std::chrono::steady_clock::now() - start).count()); - active_taskflows.fetch_sub(1, std::memory_order_relaxed); - completed_taskflows.fetch_add(1, std::memory_order_relaxed); - publish_state(); - return elapsed; + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + return run_synchronous(taskflow, nullptr); } void run(Task_Graph& taskflow, std::function completion) { if (!completion) throw std::invalid_argument("Taskflow completion is empty"); - ensure_executor(); - auto active = active_taskflows.fetch_add(1, std::memory_order_relaxed) + 1; - auto peak = peak_active_taskflows.load(std::memory_order_relaxed); - while (peak < active && !peak_active_taskflows.compare_exchange_weak( - peak, active, std::memory_order_relaxed)) {} - executor->run(native_taskflow(taskflow), [this, completion = std::move(completion)]() mutable { - active_taskflows.fetch_sub(1, std::memory_order_relaxed); - completed_taskflows.fetch_add(1, std::memory_order_relaxed); - publish_state(); + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + run_asynchronous(taskflow, false, [completion = std::move(completion)]( + detail::Taskflow_Graph_Timing) mutable { completion(); }); } @@ -954,35 +955,51 @@ public: std::string_view stage) { const auto token = detail::Taskflow_Frame_Access::begin_graph( frame, taskflow, stage); - observer->activate_graph(frame, taskflow); + /* + * 真无锁 Node::data -> Frame 路由的成立条件: + * 1. 同一 DAG 的同一批 Node 不能同时绑定不同帧并发执行;_data 只有一个普通 void* 槽位。 + * 2. Frame 必须覆盖整个 DAG 执行以及全部 Observer on_entry/on_exit 回调的生命周期。 + * 3. 多 worker 会并发附加 trace;Frame 必须按 worker 隔离单写者,不能共享写同一容器。 + */ + detail::Task_Graph_Access::bind_frame( + taskflow, frame.taskflow_trace_requested() ? &frame : nullptr); + detail::Taskflow_Graph_Timing timing{}; try { - const auto elapsed = run(taskflow); - observer->deactivate_graph(frame, taskflow); - detail::Taskflow_Frame_Access::finish_graph(token); + const auto elapsed = run_synchronous(taskflow, &timing); + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + timing.observer_unbind_finished = std::chrono::steady_clock::now(); + detail::Taskflow_Frame_Access::finish_graph(token, timing); return elapsed; } catch (...) { - observer->deactivate_graph(frame, taskflow); - detail::Taskflow_Frame_Access::finish_graph(token); + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + timing.observer_unbind_finished = std::chrono::steady_clock::now(); + detail::Taskflow_Frame_Access::finish_graph(token, timing); throw; } } void run(Task_Graph& taskflow, Render_Frame& frame, std::string_view stage, std::function completion) { + if (!completion) throw std::invalid_argument("Taskflow completion is empty"); const auto token = detail::Taskflow_Frame_Access::begin_graph( frame, taskflow, stage); try { - observer->activate_graph(frame, taskflow); - run(taskflow, [this, &taskflow, &frame, token, - completion = std::move(completion)]() mutable { - observer->deactivate_graph(frame, taskflow); - detail::Taskflow_Frame_Access::finish_graph(token); + detail::Task_Graph_Access::bind_frame( + taskflow, frame.taskflow_trace_requested() ? &frame : nullptr); + run_asynchronous(taskflow, true, [this, &taskflow, &frame, token, + completion = std::move(completion)]( + detail::Taskflow_Graph_Timing timing) mutable { + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + timing.observer_unbind_finished = std::chrono::steady_clock::now(); + detail::Taskflow_Frame_Access::finish_graph(token, timing); completion(); }); } catch (...) { - observer->deactivate_graph(frame, taskflow); - detail::Taskflow_Frame_Access::finish_graph(token); + detail::Taskflow_Graph_Timing timing{}; + detail::Task_Graph_Access::bind_frame(taskflow, nullptr); + timing.observer_unbind_finished = std::chrono::steady_clock::now(); + detail::Taskflow_Frame_Access::finish_graph(token, timing); throw; } } diff --git a/kernel/src/kernel/renderable.ipp b/kernel/src/kernel/renderable.ipp index 426a661..53c26b8 100644 --- a/kernel/src/kernel/renderable.ipp +++ b/kernel/src/kernel/renderable.ipp @@ -206,8 +206,7 @@ inline void Renderable::bind_dependency_graph_object(Attached auto* object) { [](Root* root) { auto* value = static_cast(root); auto& private_data = static_cast(*value->d); - private_data.state.advance(); - value->template notify_state(); + private_data.template publish_state(); } }, business_name diff --git a/kernel/src/kernel/scene.ipp b/kernel/src/kernel/scene.ipp index 354764e..d156b47 100644 --- a/kernel/src/kernel/scene.ipp +++ b/kernel/src/kernel/scene.ipp @@ -73,9 +73,12 @@ void Scene::Private::process(Object* object, Render_Frame* frame, Callback&& cal ? detail::run_taskflow(*runtime->taskflow, *frame, "scene.prepare") : detail::run_taskflow(*runtime->taskflow); if (frame) frame->mark(Frame_Trace_Marker::prepare_finished); + object->template current_dependency_graph().for_each_bound( + [](Renderable* renderable, Renderable::Private& data) { + data.dispatch->state.publish(renderable); + }); state.event_statistics = event_statistics.state(); - private_data.state.advance(); - object->template notify_state(); + private_data.template publish_state(); Result result; std::invoke(std::forward(callback), std::as_const(result)); } @@ -178,7 +181,6 @@ void Scene::Private::after_advance(Object* object, std::chrono::steady_clock::now().time_since_epoch()).count()); state.prepare_execution_time_ns = finished - state.prepare_execution_time_ns; } - dispatch->state.publish(root); }); prepare_if.precede(prepare_run); prepare_if.precede(prepare_done); diff --git a/kernel/src/test/object_test.cpp b/kernel/src/test/object_test.cpp index 57ae200..eb0430d 100644 --- a/kernel/src/test/object_test.cpp +++ b/kernel/src/test/object_test.cpp @@ -74,7 +74,7 @@ TEST(object_builder, build_attaches_private_and_commits_staged_values) { ASSERT_TRUE(result.has_value()); auto object = std::move(result).value(); EXPECT_EQ(object->data_for_test().current->first, 29); - EXPECT_EQ(object->data_for_test().pending->first, 0); + EXPECT_EQ(object->data_for_test().pending->first, 29); EXPECT_EQ(object->current_buffer(), 31); EXPECT_EQ(object->pending_buffer(), 31); } @@ -89,11 +89,11 @@ TEST(object_buffer, state_commits_and_keeps_incremental_baseline) { EXPECT_EQ(object->data_for_test().state.current->first, 3); EXPECT_EQ(object->data_for_test().state.current->second, 5); } -TEST(object_buffer, prop_publishes_and_keeps_incremental_baseline) { +TEST(object_buffer, prop_commits_and_keeps_incremental_baseline) { auto object = build_object(); object->set<&Test_Object::Prop::first>(17); - EXPECT_EQ(object->data_for_test().current->first, 17); - EXPECT_EQ(object->data_for_test().pending->first, 0); + EXPECT_EQ(object->data_for_test().current->first, 0); + EXPECT_EQ(object->data_for_test().pending->first, 17); object->advance(); EXPECT_EQ(object->data_for_test().pending->first, 17); EXPECT_EQ(object->data_for_test().current->first, 17); @@ -168,12 +168,11 @@ TEST(state_tag, callback_publishes_only_requested_layer) { [&](const auto& state) { ++calls; EXPECT_EQ(state.first, 23); + EXPECT_EQ(object->read_state().first, 23); } ); object->update_state<&Test_Object::State::first>(23); - object->advance(); - EXPECT_EQ(calls, 0); - object->notify_state(); + object->publish_state(); EXPECT_EQ(calls, 1); } static_assert(requires(Object& object) { @@ -191,10 +190,10 @@ TEST(state_tag, inherited_tags_remain_independently_addressable) { int derived_calls = 0; object->set_state_callback([&](const auto&) { ++base_calls; }); object->set_state_callback([&](const auto&) { ++derived_calls; }); - object->notify_state(); + object->publish_state(); EXPECT_EQ(base_calls, 1); EXPECT_EQ(derived_calls, 0); - object->notify_state(); + object->publish_state(); EXPECT_EQ(base_calls, 1); EXPECT_EQ(derived_calls, 1); } diff --git a/kernel/src/test/render_test.cpp b/kernel/src/test/render_test.cpp index f4cd236..56ca012 100644 --- a/kernel/src/test/render_test.cpp +++ b/kernel/src/test/render_test.cpp @@ -237,7 +237,7 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) aethera::detail::finish_taskflow_trace(frame); EXPECT_EQ(completed.load(), 3); - const auto trace = frame.taskflow_trace(); + const auto trace = frame.take_taskflow_trace(); ASSERT_EQ(trace.graphs.size(), 1u); EXPECT_EQ(trace.identity, (aethera::Frame_Identity{41, 73})); EXPECT_TRUE(std::ranges::contains( @@ -246,7 +246,12 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) &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(), 4u); + EXPECT_EQ(std::ranges::count_if(trace.graphs.front().nodes, [](const auto& node) { + return node.type != "diagnostic"; + }), 4); + EXPECT_EQ(std::ranges::count_if(trace.graphs.front().nodes, [](const auto& node) { + return node.type == "diagnostic"; + }), 5); const auto module = std::ranges::find( trace.graphs.front().nodes, "spectrum", &aethera::Taskflow_Graph_Trace::Node::name); @@ -283,4 +288,14 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) EXPECT_NEAR(publish_execution->ready_ms, child_paint_execution->finished_ms, 0.05); EXPECT_LT(publish_execution->queue_wait_ms, 1.0); + for (const auto name : {"taskflow.executor.finalize", "taskflow.runtime.bookkeeping", + "taskflow.runtime.publish_state", "taskflow.observer.unbind_frame", + "taskflow.graph.finish"}) { + const auto diagnostic = std::ranges::find( + trace.graphs.front().nodes, name, &aethera::Taskflow_Graph_Trace::Node::name); + ASSERT_NE(diagnostic, trace.graphs.front().nodes.end()); + EXPECT_NE(std::ranges::find(trace.tasks, diagnostic->native_id, + &aethera::Taskflow_Task_Trace::native_id), + trace.tasks.end()); + } } diff --git a/render_2D/render_2D/base/Renderable_2D.ipp b/render_2D/render_2D/base/Renderable_2D.ipp index f75db0a..ca70438 100644 --- a/render_2D/render_2D/base/Renderable_2D.ipp +++ b/render_2D/render_2D/base/Renderable_2D.ipp @@ -48,8 +48,11 @@ void Renderable_2D::Private::bind_private_crtp(Object* object) { return &static_cast(root)->template pending_buffer(); }; this->cache_enabled = [](const Root* root) { - return static_cast(root) - ->template read_prop().cache_enabled; + const auto* value = static_cast(root); + const auto& private_data = static_cast(*value->d); + const Renderable_2D::Prop& prop = + private_data.template current_prop(); + return prop.cache_enabled; }; } } diff --git a/render_2D/render_2D/plottable/Afterglow.ipp b/render_2D/render_2D/plottable/Afterglow.ipp index bc11767..725d12a 100644 --- a/render_2D/render_2D/plottable/Afterglow.ipp +++ b/render_2D/render_2D/plottable/Afterglow.ipp @@ -133,7 +133,7 @@ Task_Graph Afterglow::Private::build_paint_graph(Object* object, const Prop&) { } template void Afterglow::Private::prepare_frame(Object* object) { - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); object->template exchange_stream(); prepared.source_spectra.clear(); object->template access_rendering_stream( @@ -141,12 +141,12 @@ void Afterglow::Private::prepare_frame(Object* object) { prepared.source_spectra.assign(spectra.begin(), spectra.end()); }); object->template accumulate_stream(std::size_t{64}); - const auto& frequency_layout = frequency_axis->template read_prop(); - const auto& power_layout = power_axis->template read_prop(); + const auto& frequency_layout = this->template read_current_prop(frequency_axis); + const auto& power_layout = this->template read_current_prop(power_axis); const std::size_t available = prepared.source_spectra.empty() ? 0 : prepared.source_spectra.back()->size(); const int columns = static_cast(state.frequency_point_size ? std::min(state.frequency_point_size, available) : available); const int rows = static_cast(state.power_point_size ? state.power_point_size : std::max(1.0, std::abs(power_layout.pixel_length))); - const Size canvas = scene->template read_prop().viewport; + const Size canvas = this->template read_current_prop(scene).viewport; const auto layout = detail::raster_layout(frequency_axis, state.frequency_range, columns, power_axis, state.power_range, rows, frequency_layout.orientation, power_layout.orientation); if (canvas.empty() || !layout.valid()) { prepared.valid = false; @@ -170,7 +170,7 @@ void Afterglow::Private::prepare_frame(Object* object) { template void Afterglow::Private::accumulate_partition(Object* object, Plot_Partition_Count index) { if (!prepared.valid) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const int columns = prepared.layout.first_horizontal ? prepared.layout.width : prepared.layout.height; const int rows = prepared.layout.first_horizontal ? prepared.layout.height : prepared.layout.width; const detail::Raster_Partition partition = detail::raster_partition( @@ -219,7 +219,7 @@ inline void Afterglow::Private::normalize_frame() { template void Afterglow::Private::color_partition(Object* object, Plot_Partition_Count index) { if (!prepared.valid) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const int columns = prepared.layout.first_horizontal ? prepared.layout.width : prepared.layout.height; const detail::Raster_Partition partition = detail::raster_partition( prepared.layout, index, graph_partition_grid); diff --git a/render_2D/render_2D/plottable/Constellation_Diagram.ipp b/render_2D/render_2D/plottable/Constellation_Diagram.ipp index e9c35bb..a2c780f 100644 --- a/render_2D/render_2D/plottable/Constellation_Diagram.ipp +++ b/render_2D/render_2D/plottable/Constellation_Diagram.ipp @@ -21,10 +21,10 @@ struct Constellation_Diagram::Private : Prev_Private { /* CRTP 覆盖:构建锚点前驱与并行接收点 Paint 子图。 */ template [[nodiscard]] Task_Graph build_paint_graph(Object* object, - const Prop& state); + const Prop& state); template [[nodiscard]] bool should_rebuild_paint_graph(Object* object, - const Prop& state); + const Prop& state); template void paint_anchors(Object* object); template @@ -38,8 +38,7 @@ struct Constellation_Diagram::Private : Prev_Private { template Constellation_Diagram::Builder::Builder( Axis_Object* i_axis_value, Axis_Object* q_axis_value, - Plot_Partition_Count partition_count) - : Base(), i_axis(i_axis_value), q_axis(q_axis_value) { + Plot_Partition_Count partition_count) : Base(), i_axis(i_axis_value), q_axis(q_axis_value) { this->prop.partition_count = partition_count; } template @@ -59,8 +58,8 @@ std::expected, Dependency_Graph_Error> Constellation_Dia prepare.add_dependency(object, scene); prepare.add_dependency(object, i_axis); prepare.add_dependency(object, q_axis); - paint.add_dependency(i_axis, object); - paint.add_dependency(q_axis, object); + // paint.add_dependency(i_axis, object); + // paint.add_dependency(q_axis, object); cache.template add_prop_dependency<&Render_Scene_2D::Prop::viewport>(object, scene); cache.template add_prop_dependency<&Abs_Axis::Prop::position>(object, i_axis); cache.template add_prop_dependency<&Abs_Axis::Prop::pixel_length>(object, i_axis); @@ -76,17 +75,17 @@ std::expected, Dependency_Graph_Error> Constellation_Dia } template void Constellation_Diagram::Private::prepare_data(Object* object) { - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); object->template exchange_stream(); const auto current = monotonic_milliseconds(); object->template accumulate_stream( [&](const Constellation_Point& point) { return current - point.submitted_at_ms <= state.point_lifetime_ms; }); - const auto& i_layout = i_axis->template read_prop(); - const auto& q_layout = q_axis->template read_prop(); + const auto& i_layout = this->template read_current_prop(i_axis); + const auto& q_layout = this->template read_current_prop(q_axis); prepared = {}; - prepared.canvas = scene->template read_prop().viewport; + prepared.canvas = this->template read_current_prop(scene).viewport; if (prepared.canvas.empty() || i_layout.orientation == q_layout.orientation) { return; } @@ -144,16 +143,18 @@ inline Rect_F Constellation_Diagram::Private::point_region( } constexpr double antialias_padding{1.0}; const double padding = radius + antialias_padding; - return {left - padding, top - padding, - right - left + padding * 2.0, - bottom - top + padding * 2.0}; + return { + left - padding, top - padding, + right - left + padding * 2.0, + bottom - top + padding * 2.0 + }; } template void Constellation_Diagram::Private::paint_anchors(Object* object) { - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); if (!prepared.valid || prepared.anchors.empty()) return; detail::Painter painter(this->paint_surface(), prepared.canvas, - point_region(prepared.anchors, 4.5)); + point_region(prepared.anchors, 4.5)); /* * 原实现每个点都会创建 BLPath,再分别 fill + stroke;这里改为同色实心圆批量入口。 * 旧样式是 radius + 1px 同色描边,最终外半径约增加 0.5px,因此直接使用 4.5/2.5 @@ -165,23 +166,25 @@ template void Constellation_Diagram::Private::paint_partition( Object* object, Plot_Partition_Count index) { const auto& state = - object->template read_prop(); + this->template read_current_prop(object); if (!prepared.valid || state.partition_count == 0 || index >= state.partition_count) return; const std::size_t point_count = prepared.points.size(); const auto boundary = [point_count, count = state.partition_count]( - Plot_Partition_Count position) { + Plot_Partition_Count position) { return point_count / count * position + - std::min(position, point_count % count); + std::min(position, point_count % count); }; const std::size_t first = boundary(index); const std::size_t last = boundary(index + 1); - const std::span points{prepared.points.data() + first, - last - first}; + const std::span points{ + prepared.points.data() + first, + last - first + }; if (points.empty()) return; detail::Painter painter(this->paint_surface(), prepared.canvas, - point_region(points, 2.5)); + point_region(points, 2.5)); painter.solid_circles(points, 2.5, state.point_color); } template diff --git a/render_2D/render_2D/plottable/Frequency_Trace.ipp b/render_2D/render_2D/plottable/Frequency_Trace.ipp index 2b5c21b..097cbe3 100644 --- a/render_2D/render_2D/plottable/Frequency_Trace.ipp +++ b/render_2D/render_2D/plottable/Frequency_Trace.ipp @@ -59,11 +59,11 @@ Task_Graph Frequency_Trace::Private::build_paint_graph(Object* object, const Pro } template void Frequency_Trace::Private::prepare_frame(Object* object, Plot_Partition_Count partition_count) { - const auto& state = object->template read_prop(); const auto& time_layout = time_axis->template read_prop(); const auto& value_layout = value_axis->template read_prop(); - const auto visible_count = static_cast(std::max(2, time_axis->template read_prop().visible_count)); + const auto& state = this->template read_current_prop(object); const auto& time_layout = this->template read_current_prop(time_axis); const auto& value_layout = this->template read_current_prop(value_axis); + const auto visible_count = static_cast(std::max(2, this->template read_current_prop(time_axis).visible_count)); object->template exchange_stream(); object->template accumulate_stream(visible_count); - prepared = {}; prepared.partitions.resize(partition_count); prepared.canvas = scene->template read_prop().viewport; + prepared = {}; prepared.partitions.resize(partition_count); prepared.canvas = this->template read_current_prop(scene).viewport; if (prepared.canvas.empty() || time_layout.orientation == value_layout.orientation) return; object->template access_rendering_stream([&](std::span samples) { prepared.samples.reserve(samples.size()); @@ -74,14 +74,14 @@ void Frequency_Trace::Private::prepare_frame(Object* object, Plot_Partition_Coun } template void Frequency_Trace::Private::prepare_partition(Object* object, Plot_Partition_Count partition_index) { - if (!prepared.valid) return; const auto& state = object->template read_prop(); const auto& time_layout = time_axis->template read_prop(); const auto& value_layout = value_axis->template read_prop(); const auto& value_state = value_axis->template read_prop(); + if (!prepared.valid) return; const auto& state = this->template read_current_prop(object); const auto& time_layout = this->template read_current_prop(time_axis); const auto& value_layout = this->template read_current_prop(value_axis); const auto& value_state = this->template read_current_prop(value_axis); const Axis_Range domain{prepared.samples.front().coordinate, prepared.samples.back().coordinate}; const auto range = detail::curve_partition_range(prepared.samples.size(), partition_index, prepared.partitions.size(), domain); prepared.partitions[partition_index] = detail::prepare_curve(std::span(prepared.samples).subspan(range.first_sample, range.sample_count), true, time_axis->coordinate_range(), value_state.coordinate_range, time_axis, value_axis, time_layout.orientation, value_layout.orientation); } template void Frequency_Trace::Private::paint_partition(Object* object, Plot_Partition_Count partition_index) { if (!prepared.valid || partition_index >= prepared.partitions.size()) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const auto& curve = prepared.partitions[partition_index]; const Rect_F region = detail::curve_paint_region(curve, state.pen.width + 2.0); if (region.empty()) return; diff --git a/render_2D/render_2D/plottable/Selection_Rectangle_Overlay.ipp b/render_2D/render_2D/plottable/Selection_Rectangle_Overlay.ipp index 5f5a257..b0f52b7 100644 --- a/render_2D/render_2D/plottable/Selection_Rectangle_Overlay.ipp +++ b/render_2D/render_2D/plottable/Selection_Rectangle_Overlay.ipp @@ -59,7 +59,7 @@ template void Selection_Rectangle_Overlay::Private::paint(Object* object) { auto& data = static_cast(*this); const auto& state = static_cast(*data.current); - const Size canvas = scene->template read_prop().viewport; + const Size canvas = this->template read_current_prop(scene).viewport; auto& cache = data.paint_surface(); if (canvas.empty()) return; detail::Painter painter(cache, canvas); diff --git a/render_2D/render_2D/plottable/Spectrum.ipp b/render_2D/render_2D/plottable/Spectrum.ipp index 1db0208..6b97031 100644 --- a/render_2D/render_2D/plottable/Spectrum.ipp +++ b/render_2D/render_2D/plottable/Spectrum.ipp @@ -163,10 +163,10 @@ void Spectrum::Private::prepare_frame(Object* object, std::size_t partition_coun minima[index] = std::min(minima[index], frame.samples[index]); } } - const auto& scene_state = scene->template read_prop(); - const auto& frequency_layout = frequency_axis->template read_prop(); - const auto& power_layout = power_axis->template read_prop(); - const auto& power_state = power_axis->template read_prop(); + const auto& scene_state = this->template read_current_prop(scene); + const auto& frequency_layout = this->template read_current_prop(frequency_axis); + const auto& power_layout = this->template read_current_prop(power_axis); + const auto& power_state = this->template read_current_prop(power_axis); prepared = {}; prepared.partitions.resize(partition_count); if (frequency_layout.orientation == power_layout.orientation) return; @@ -199,10 +199,10 @@ void Spectrum::Private::prepare_partition(Object* object, std::size_t partition_ const auto& state = static_cast(*private_data.current); const auto& frame = object->template current_buffer(); if (frame.samples.empty()) return; - const auto& frequency_layout = frequency_axis->template read_prop(); - const auto& power_layout = power_axis->template read_prop(); - const auto& frequency_state = frequency_axis->template read_prop(); - const auto& power_state = power_axis->template read_prop(); + const auto& frequency_layout = this->template read_current_prop(frequency_axis); + const auto& power_layout = this->template read_current_prop(power_axis); + const auto& frequency_state = this->template read_current_prop(frequency_axis); + const auto& power_state = this->template read_current_prop(power_axis); const auto range = detail::curve_partition_range(frame.samples.size(), partition_index, prepared.partitions.size(), state.frequency_range); auto& partition = prepared.partitions[partition_index]; partition.clip = detail::map_plot_rect(frequency_axis, range.domain, power_axis, power_state.coordinate_range, frequency_layout.orientation); diff --git a/render_2D/render_2D/plottable/Sweep_Spectrum.ipp b/render_2D/render_2D/plottable/Sweep_Spectrum.ipp index 5d521bf..74e96e9 100644 --- a/render_2D/render_2D/plottable/Sweep_Spectrum.ipp +++ b/render_2D/render_2D/plottable/Sweep_Spectrum.ipp @@ -72,7 +72,7 @@ Task_Graph Sweep_Spectrum::Private::build_paint_graph(Object* object, const Prop } template void Sweep_Spectrum::Private::prepare_frame(Object* object) { - const auto& state = object->template read_prop(); const auto& frequency_layout = frequency_axis->template read_prop(); const auto& power_layout = power_axis->template read_prop(); const auto& power_state = power_axis->template read_prop(); prepared = {}; prepared.canvas = scene->template read_prop().viewport; + const auto& state = this->template read_current_prop(object); const auto& frequency_layout = this->template read_current_prop(frequency_axis); const auto& power_layout = this->template read_current_prop(power_axis); const auto& power_state = this->template read_current_prop(power_axis); prepared = {}; prepared.canvas = this->template read_current_prop(scene).viewport; const std::size_t block_count = std::max(1, state.block_count); if (sweep_blocks.size() != block_count) { sweep_blocks.assign(block_count, {}); @@ -109,13 +109,13 @@ void Sweep_Spectrum::Private::prepare_frame(Object* object) { } template void Sweep_Spectrum::Private::prepare_partition(Object* object, Plot_Partition_Count index) { - if (!prepared.valid) return; const auto& state = object->template read_prop(); const auto& frequency_layout = frequency_axis->template read_prop(); const auto& power_layout = power_axis->template read_prop(); const auto& frequency_state = frequency_axis->template read_prop(); const auto& power_state = power_axis->template read_prop(); const auto range = detail::curve_partition_range(prepared.values.size(), index, prepared.partitions.size(), prepared.domain); + if (!prepared.valid) return; const auto& state = this->template read_current_prop(object); const auto& frequency_layout = this->template read_current_prop(frequency_axis); const auto& power_layout = this->template read_current_prop(power_axis); const auto& frequency_state = this->template read_current_prop(frequency_axis); const auto& power_state = this->template read_current_prop(power_axis); const auto range = detail::curve_partition_range(prepared.values.size(), index, prepared.partitions.size(), prepared.domain); prepared.partitions[index] = detail::prepare_curve(std::span(prepared.values).subspan(range.first_sample, range.sample_count), range.domain, state.interpolation_mode, state.visible_range_only, frequency_state.coordinate_range, power_state.coordinate_range, frequency_axis, power_axis, frequency_layout.orientation, power_layout.orientation); } template void Sweep_Spectrum::Private::paint_partition(Object* object, Plot_Partition_Count index) { if (!prepared.valid || index >= prepared.partitions.size()) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const auto& curve = prepared.partitions[index]; const Rect_F region = detail::curve_paint_region(curve, state.pen.width + 2.0); if (region.empty()) return; @@ -125,7 +125,7 @@ void Sweep_Spectrum::Private::paint_partition(Object* object, Plot_Partition_Cou template void Sweep_Spectrum::Private::paint_marker(Object* object) { if (!prepared.valid) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const double padding = state.current_frequency_pen.width + 2.0; const auto [left, right] = std::minmax(prepared.marker_first.x, prepared.marker_second.x); const auto [top, bottom] = std::minmax(prepared.marker_first.y, prepared.marker_second.y); diff --git a/render_2D/render_2D/plottable/Waterfall.ipp b/render_2D/render_2D/plottable/Waterfall.ipp index ebadfae..bfa748a 100644 --- a/render_2D/render_2D/plottable/Waterfall.ipp +++ b/render_2D/render_2D/plottable/Waterfall.ipp @@ -138,10 +138,10 @@ Task_Graph Waterfall::Private::build_paint_graph(Object* object, const Prop&) { } template void Waterfall::Private::prepare_frame(Object* object) { - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); object->template exchange_stream(); const auto row_limit = static_cast(std::max( - 2, time_axis->template read_prop().visible_count)); + 2, this->template read_current_prop(time_axis).visible_count)); object->template accumulate_stream(row_limit); prepared = {}; std::vector> source_rows; @@ -149,9 +149,9 @@ void Waterfall::Private::prepare_frame(Object* object) { [&](std::span> rows) { source_rows.assign(rows.begin(), rows.end()); }); - const auto& frequency_layout = frequency_axis->template read_prop(); - const auto& time_layout = time_axis->template read_prop(); - prepared.canvas = scene->template read_prop().viewport; + const auto& frequency_layout = this->template read_current_prop(frequency_axis); + const auto& time_layout = this->template read_current_prop(time_axis); + prepared.canvas = this->template read_current_prop(scene).viewport; if (source_rows.empty()) { return; } @@ -164,7 +164,7 @@ void Waterfall::Private::prepare_frame(Object* object) { if (!selection) { return; } - const int time_slots = std::max(2, time_axis->template read_prop().visible_count); + const int time_slots = std::max(2, this->template read_current_prop(time_axis).visible_count); const Axis_Range time_range = time_axis->coordinate_range(); prepared.layout = detail::raster_layout(frequency_axis, selection->range, selection->count(), time_axis, time_range, time_slots, frequency_layout.orientation, time_layout.orientation); if (prepared.canvas.empty() || !prepared.layout.valid()) { @@ -199,7 +199,7 @@ void Waterfall::Private::prepare_frame(Object* object) { template void Waterfall::Private::prepare_partition(Object* object, Plot_Partition_Count index) { if (!prepared.valid) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); const detail::Raster_Partition partition = detail::raster_partition( prepared.layout, index, graph_partition_grid); for (int y = partition.first_y; y < partition.last_y; ++y) { @@ -225,7 +225,7 @@ void Waterfall::Private::prepare_partition(Object* object, Plot_Partition_Count template void Waterfall::Private::paint_partition(Object* object, Plot_Partition_Count index) { - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); if (!prepared.valid || index >= detail::raster_partition_count(graph_partition_grid)) return; @@ -238,7 +238,7 @@ void Waterfall::Private::paint_partition(Object* object, template void Waterfall::Private::paint_tooltip(Object* object) { if (!prepared.valid || prepared.tooltip_text.empty()) return; - const auto& state = object->template read_prop(); + const auto& state = this->template read_current_prop(object); detail::Painter painter(this->paint_surface(), prepared.canvas, prepared.tooltip_box); painter.rect(prepared.tooltip_box, Pen{state.tooltip_text_pen.color}, diff --git a/render_2D/render_2D/scene/Render_Scene_2D.hpp b/render_2D/render_2D/scene/Render_Scene_2D.hpp index 814c8e0..f23edba 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.hpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.hpp @@ -29,20 +29,14 @@ struct Render_Scene_2D : Def; using Final_Builder = typename Object::Builder; Builder(); - template requires std::derived_from < Renderable_Object - - , - Renderable_2D_Base - > + template requires std::derived_from Final_Builder& add_renderable(Renderable_Object* renderable); [[nodiscard]] std::expected, Dependency_Graph_Error> build(); private: - std::vector(Object *)> - > - attachments {}; /* 仅在 build 期间绑定已构造 Renderable。 */ + std::vector(Object*)>> attachments{}; /* 仅在 build 期间绑定已构造 Renderable。 */ }; enum struct Render_Result { frame_in_flight, view_inactive, empty_viewport }; - using Frame_Callback = std::function; + using Frame_Callback = std::function; /* 向调用方拥有的帧合成一次;Scene::advance 到完成像素发布保持不可重入。 */ [[nodiscard]] std::expected render(Frame_2D* frame); /* diff --git a/render_2D/render_2D/scene/Render_Scene_2D.ipp b/render_2D/render_2D/scene/Render_Scene_2D.ipp index 9abbf33..ecd5749 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.ipp +++ b/render_2D/render_2D/scene/Render_Scene_2D.ipp @@ -200,7 +200,6 @@ void Render_Scene_2D::Private::after_advance(Object* object, Prop*, State_Access std::chrono::steady_clock::now().time_since_epoch()).count()); state.paint_execution_time_ns = finished - state.paint_execution_time_ns; } - dispatch->state.publish(root); }); paint_if.precede(paint_run); paint_if.precede(paint_done); @@ -377,7 +376,8 @@ void Render_Scene_2D::Private::ensure_frame_taskflow(Object* object) { 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(); + const auto& current_data = static_cast(*this); + const Prop& prop = current_data.template current_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); @@ -429,7 +429,8 @@ void Render_Scene_2D::Private::ensure_frame_taskflow(Object* object) { auto* frame_object = active_frame; if (!frame_object || !frame_target) throw std::logic_error("2D frame target preparation lost its frame"); - const auto& prop = object->template read_prop(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_prop(); frame_object->mark(Frame_Trace_Marker::paint_frame_target_started); frame_target->ensure_size(prop.viewport); frame_target->clear(); @@ -442,7 +443,8 @@ void Render_Scene_2D::Private::ensure_frame_taskflow(Object* object) { auto* frame_object = active_frame; if (!frame_object || !frame_target) throw std::logic_error("2D background paint lost its frame target"); - const auto& prop = object->template read_prop(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_prop(); frame_object->mark(Frame_Trace_Marker::paint_background_started); { detail::Painter painter(*frame_target, prop.viewport); @@ -460,7 +462,8 @@ void Render_Scene_2D::Private::ensure_frame_taskflow(Object* object) { auto* frame_object = active_frame; if (!frame_object || !frame_target) throw std::logic_error("2D cache target preparation lost its frame"); - const auto& prop = object->template read_prop(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_prop(); frame_object->mark(Frame_Trace_Marker::paint_cache_targets_started); prepare_paint_targets(object, *frame_target, prop.viewport); frame_object->mark(Frame_Trace_Marker::paint_cache_targets_finished); @@ -525,8 +528,8 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { * results; it must not modify a graph that is still executing. */ object->advance(); - const auto& prop = - object->template read_prop(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_prop(); if (!prop.view_active) { release_admission(); return std::unexpected(Render_Result::view_inactive); @@ -555,11 +558,22 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { * corun 外接 H264/WebRTC DAG,而下一帧已经允许开始。 */ auto& private_data = static_cast(*this); + std::unordered_set published_renderables; + const auto publish_renderable_states = [&](const auto& graph) { + graph.for_each_bound([&](Renderable* renderable, + Renderable::Private& data) { + if (published_renderables.insert(renderable).second) + data.dispatch->state.publish(renderable); + }); + }; + publish_renderable_states( + object->template current_dependency_graph()); + publish_renderable_states( + object->template current_dependency_graph()); 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(); + private_data.template publish_state(); active_frame = nullptr; frame_target = nullptr; diff --git a/render_3D/render_3D/scene/Render_Scene_3D.ipp b/render_3D/render_3D/scene/Render_Scene_3D.ipp index 06e5125..9d7f258 100644 --- a/render_3D/render_3D/scene/Render_Scene_3D.ipp +++ b/render_3D/render_3D/scene/Render_Scene_3D.ipp @@ -108,7 +108,9 @@ Render_Scene_3D::Builder::add_camera(Camera_Object* camera_value) { if (camera) throw std::logic_error("Render_Scene_3D accepts one Camera component"); camera = camera_value; read_camera = [](const Root* root) { - const auto& prop = static_cast(root)->template read_prop(); + const auto* object = static_cast(root); + const Camera_3D::Prop& prop = + Base::template current_prop(object); return Camera_Descriptor{prop.initial_view, prop.projection, prop.controller, prop.turntable_control, prop.arcball_control, prop.fly_control, prop.panzoom_control, @@ -125,7 +127,9 @@ Render_Scene_3D::Builder::add_axes(Axes_Object* axes_value) { if (axes) throw std::logic_error("Render_Scene_3D accepts one Axes component"); axes = axes_value; read_axes = [](const Root* root) { - const auto& prop = static_cast(root)->template read_prop(); + const auto* object = static_cast(root); + const Axes_3D::Prop& prop = + Base::template current_prop(object); return std::array{prop.x_axis, prop.y_axis, prop.z_axis}; }; return static_cast(*this); @@ -158,7 +162,8 @@ std::expected, Dependency_Graph_Error> Render_Scene_3D:: } template detail::Scene_3D_Parameters Render_Scene_3D::Private::parameters(Object* object) const { - const auto& prop = object->template read_prop(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_prop(); const auto axis = read_axes(axes_component); return {prop.viewport, prop.clear_color, read_camera(camera_component), axis[0], axis[1], axis[2]}; } @@ -286,7 +291,6 @@ void Render_Scene_3D::Private::ensure_frame_taskflow(Object* object) { std::chrono::steady_clock::now().time_since_epoch()).count()); state.paint_execution_time_ns = finished - state.paint_execution_time_ns; } - dispatch->state.publish(root); }); done.describe("renderable", std::string(dispatch->business_name)) .describe("dimension", "3D") @@ -390,7 +394,6 @@ void Render_Scene_3D::Private::ensure_frame_taskflow(Object* object) { std::chrono::steady_clock::now() - started) .count()); } - dispatch->state.publish(root); }); task.describe("renderable", std::string(dispatch->business_name)) .describe("dimension", "3D") @@ -491,7 +494,8 @@ void Render_Scene_3D::Private::ensure_frame_taskflow(Object* object) { 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(); + const auto& private_data = static_cast(*this); + const Prop& prop = private_data.template current_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; @@ -509,8 +513,19 @@ Render_Scene_3D::Render_Result Render_Scene_3D::Private::render(Object* object, try { auto completion = [this, object, frame] { auto& private_data = static_cast(*this); - private_data.state.advance(); - object->template notify_state(); + std::unordered_set published_renderables; + const auto publish_renderable_states = [&](const auto& graph) { + graph.for_each_bound([&](Renderable* renderable, + Renderable::Private& data) { + if (published_renderables.insert(renderable).second) + data.dispatch->state.publish(renderable); + }); + }; + publish_renderable_states( + object->template current_dependency_graph()); + publish_renderable_states( + object->template current_dependency_graph()); + private_data.template publish_state(); complete_frame(object, frame); }; if (frame->taskflow_trace_requested()) diff --git a/render_3D/render_3D/visual/Basic_Visual.ipp b/render_3D/render_3D/visual/Basic_Visual.ipp index acb8f95..e53c647 100644 --- a/render_3D/render_3D/visual/Basic_Visual.ipp +++ b/render_3D/render_3D/visual/Basic_Visual.ipp @@ -46,7 +46,7 @@ template template void Basic_Visual::Private::prepare_data(Object* object) { using Visual_Tag = typename Basic_Visual::Base_Tag; - const auto& prop = object->template read_prop(); + const auto& prop = this->template read_current_prop(object); if (!valid_items(prop.items) || !detail::finite(prop.transform)) throw std::invalid_argument("3D visual property contains invalid data"); auto data = std::make_shared(); Spec::prepare(prop.items, *data); diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index a8968dc..f2ab3e0 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -921,7 +921,7 @@ void Plot::Private::consume_completed_frame(Render_Frame* frame) { owner->d->store_trace( owner->d->post_publish_trace_control, owner->d->post_publish_trace_slots, - trace_frame->taskflow_trace()); + trace_frame->take_taskflow_trace()); } owner->d->post_publish_busy.store(false, std::memory_order_release); }; @@ -969,7 +969,7 @@ void Plot::Private::retire_completed_frame(Render_Frame* frame) { if (frame->taskflow_trace_requested()) store_trace(taskflow_trace_control, taskflow_trace_slots, - frame->taskflow_trace()); + frame->take_taskflow_trace()); auto expected = Frame_State::consuming; if (!managed->state.compare_exchange_strong( diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index 1d6b176..43e1bc8 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -1350,6 +1350,7 @@ function task_cpu_label(sample: Taskflow_Execution_Trace) { type Taskflow_Node_Phase = "prepare" | "paint" | "control" | "completion" | "other"; function taskflow_node_phase(name: string): Taskflow_Node_Phase { + if (name.startsWith("taskflow.")) return "completion"; if (name.includes(".prepare.")) return "prepare"; if (name.includes(".paint.")) return "paint"; if (name === "scene.paint" || name === "scene.paint.setup" || name === "scene.paint.complete") return "paint"; @@ -1379,6 +1380,11 @@ function taskflow_node_name(name: string) { if (name === "gallery.websocket_pixels.pack") return "WebSocket · 原生像素封装"; if (name === "gallery.websocket_pixels.publish") return "WebSocket · 原生像素发布"; if (name === "gallery.sample.complete") return "Gallery · 采样完成"; + if (name === "taskflow.executor.finalize") return "Taskflow · Executor topology 收尾"; + if (name === "taskflow.runtime.bookkeeping") return "Taskflow · 运行计数更新"; + if (name === "taskflow.runtime.publish_state") return "Taskflow · 运行时状态发布"; + if (name === "taskflow.observer.deactivate_graph") return "Taskflow · Observer graph 退役"; + if (name === "taskflow.graph.finish") return "Taskflow · Graph 完成记录"; const owner = parts[0]; const stage = parts[1]; @@ -1477,13 +1483,21 @@ function Taskflow_Node_Label({node, sample, summary, state, level}: { {node.name} - {summary ? <> + {summary ? node.type === "diagnostic" ? <> + 运行时收尾 · 样本 {summary.executed}/{summary.frames} + 墙钟 {milliseconds(summary.duration.average)} ± {milliseconds(summary.duration.variability)} · P95 {milliseconds(summary.duration.p95)} + : <> {node.type} · 样本 {summary.executed}/{summary.frames} · W {summary.workers.join(", ") || "--"} 执行 {milliseconds(summary.duration.average)} ± {milliseconds(summary.duration.variability)} · P95 {milliseconds(summary.duration.p95)} 排队 {milliseconds(summary.queue.average)} ± {milliseconds(summary.queue.variability)} · P99 {milliseconds(summary.queue.p99)} 线程 CPU{summary.cpu_time_coarse ? "(低分辨率)" : ""} {milliseconds(summary.cpu.average)} · CPU 活动 {cpu_cycles(summary.cycles.average)} 协作等待 {milliseconds(summary.cooperative_wait.average)} · Observer {milliseconds(summary.observer.average)} - : <>{node.type} · W{sample?.worker_id ?? "--"} + : node.type === "diagnostic" ? <> + 运行时收尾 + 开始 {sample ? `+${sample.started_ms.toFixed(3)} ms` : "--"} · 结束 {sample ? `+${sample.finished_ms.toFixed(3)} ms` : "--"} + 墙钟 {sample ? milliseconds(duration) : "未执行"} + : <> + {node.type} · W{sample?.worker_id ?? "--"} 就绪 {sample ? `+${sample.ready_ms.toFixed(3)} ms` : "--"} · 开始 {sample ? `+${sample.started_ms.toFixed(3)} ms` : "--"} 任务体 {sample ? milliseconds(duration) : "未执行"} · 排队 {sample ? milliseconds(wait) : "--"} 线程 CPU {sample ? task_cpu_label(sample) : "--"}{sample?.cpu_time_coarse ? "(低分辨率)" : ""} · CPU 活动 {sample ? cpu_cycles(sample.cpu_cycles) : "--"} @@ -1501,36 +1515,39 @@ function taskflow_graph_analysis(graph: Taskflow_Graph_Trace, executions: Taskfl const native_ids = new Set(graph.nodes.map(node => node.native_id)); const rows = executions.filter(value => native_ids.has(value.native_id)); const node_by_native_id = new Map(graph.nodes.map(node => [node.native_id, node])); - // Taskflow 的 MODULE Observer 区间包住了整个子图。父 Module 与子节点相加会 - // 重复计算同一段时间,因此执行量只用叶子任务;Module 仅作为包络单独展示。 - const leaf_rows = rows.filter(row => + const diagnostic_rows = rows.filter(row => + node_by_native_id.get(row.native_id)?.type.toLowerCase() === "diagnostic"); + const task_rows = rows.filter(row => + node_by_native_id.get(row.native_id)?.type.toLowerCase() !== "diagnostic"); + const leaf_rows = task_rows.filter(row => node_by_native_id.get(row.native_id)?.type.toLowerCase() !== "module"); - const module_rows = rows.filter(row => + const module_rows = task_rows.filter(row => node_by_native_id.get(row.native_id)?.type.toLowerCase() === "module"); const wall_time = Math.max(0, graph.finished_ms - graph.submitted_ms); const execution_time = leaf_rows.reduce((sum, row) => sum + row.duration_ms, 0); const module_envelope_time = module_rows.reduce((sum, row) => sum + row.duration_ms, 0); + const diagnostic_time = diagnostic_rows.reduce((sum, row) => sum + row.duration_ms, 0); const cpu_execution_time = leaf_rows.reduce((sum, row) => sum + row.cpu_duration_ms, 0); const cooperative_wait_time = leaf_rows.reduce((sum, row) => sum + row.cooperative_wait_ms, 0); const cpu_cycles_total = leaf_rows.reduce((sum, row) => sum + row.cpu_cycles, 0); const cpu_time_coarse = leaf_rows.some(row => row.cpu_time_coarse); const unattributed_wall_time = Math.max(0, execution_time - cooperative_wait_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( + const observer_entry_time = task_rows.reduce((sum, row) => sum + row.observer_entry_ms, 0); + const observer_exit_time = task_rows.reduce((sum, row) => sum + row.observer_exit_ms, 0); + const observer_cpu_time = task_rows.reduce((sum, row) => sum + row.observer_entry_cpu_ms + row.observer_exit_cpu_ms, 0); + const queue_time = task_rows.reduce((sum, row) => sum + row.queue_wait_ms, 0); + const first_entered = task_rows.length ? Math.min(...task_rows.map(row => row.entered_ms)) : graph.submitted_ms; + const last_completed = task_rows.length ? Math.max(...task_rows.map(row => row.completed_ms)) : graph.finished_ms; + const longest = task_rows.reduce( (result, row) => !result || row.duration_ms > result.duration_ms ? row : result, null); - const longest_unattributed = rows.reduce((result, row) => { + const longest_unattributed = task_rows.reduce((result, row) => { const value = Math.max(0, row.duration_ms - row.cooperative_wait_ms - row.cpu_duration_ms); const previous = result ? Math.max(0, result.duration_ms - result.cooperative_wait_ms - result.cpu_duration_ms) : -1; return value > previous ? row : result; }, null); - const events = rows.flatMap(row => [ + const events = task_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); @@ -1546,21 +1563,25 @@ function taskflow_graph_analysis(graph: Taskflow_Graph_Trace, executions: Taskfl .sort((left, right) => left - right); let body_wall_time = 0; let observer_wall_time = 0; + let diagnostic_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)) + if (task_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) || + else if (task_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 if (diagnostic_rows.some(row => row.started_ms <= middle && middle < row.finished_ms)) + diagnostic_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 unattributed_completion_tail = Math.max(0, completion_tail - diagnostic_wall_time); const node_by_id = new Map(graph.nodes.map(node => [node.native_id, node])); const levels = new Map(); const visiting = new Set(); @@ -1581,15 +1602,16 @@ function taskflow_graph_analysis(graph: Taskflow_Graph_Trace, executions: Taskfl 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, leaf_rows, module_rows, levels, wall_time, execution_time, - module_envelope_time, cpu_execution_time, cooperative_wait_time, cpu_cycles_total, cpu_time_coarse, - unattributed_wall_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), + rows, task_rows, leaf_rows, module_rows, diagnostic_rows, levels, wall_time, execution_time, + module_envelope_time, diagnostic_time, diagnostic_wall_time, cpu_execution_time, cooperative_wait_time, + cpu_cycles_total, cpu_time_coarse, unattributed_wall_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, unattributed_completion_tail, + internal_idle_time: Math.max(0, idle_wall_time - initial_wait - unattributed_completion_tail), + peak_entry_queue: task_rows.reduce((maximum, row) => Math.max(maximum, row.worker_queue_size), 0), maximum_parallelism, - node_count: graph.nodes.length, + node_count: graph.nodes.filter(node => node.type.toLowerCase() !== "diagnostic").length, + diagnostic_node_count: graph.nodes.filter(node => node.type.toLowerCase() === "diagnostic").length, layer_count: width_by_level.size, parallel_layer_count: [...width_by_level.values()].filter(width => width > 1).length, longest, longest_unattributed, @@ -1608,6 +1630,8 @@ const taskflow_graph_metric_definitions = [ {key: "wall_time", label: "Topology 墙钟", unit: "milliseconds", description: "从提交 Taskflow 到整个 Topology 完成的墙钟时间。"}, {key: "execution_time", label: "叶子任务墙钟总和", unit: "milliseconds", description: "每帧所有叶子任务执行墙钟之和。"}, {key: "module_envelope_time", label: "Module 包络总和", unit: "milliseconds", description: "每帧所有 Taskflow Module Observer 包络之和。"}, + {key: "diagnostic_time", label: "收尾诊断总和", unit: "milliseconds", description: "Executor 返回后的运行计数、状态发布、Observer graph 退役与完成记录墙钟之和。"}, + {key: "diagnostic_wall_time", label: "收尾诊断并集", unit: "milliseconds", description: "Topology 尾部被明确归因到诊断收尾节点的时间区间并集。"}, {key: "cpu_execution_time", label: "叶子任务线程 CPU", unit: "milliseconds", description: "每帧叶子任务独占 Worker 的线程 CPU 时间之和。"}, {key: "cpu_cycles_total", label: "叶子任务 CPU 活动", unit: "cycles", description: "每帧叶子任务 QueryThreadCycleTime 活动量之和。"}, {key: "cooperative_wait_time", label: "累计协作等待", unit: "milliseconds", description: "每帧 Task_Graph 协作等待墙钟之和。"}, @@ -1621,7 +1645,8 @@ const taskflow_graph_metric_definitions = [ {key: "queue_time", label: "累计节点排队", unit: "milliseconds", description: "各节点从前驱完成到进入 Worker 的估算等待之和。"}, {key: "initial_wait", label: "首次准入等待", unit: "milliseconds", description: "Topology 提交后首节点进入 Worker 前的等待。"}, {key: "internal_idle_time", label: "内部调度空洞", unit: "milliseconds", description: "首尾等待之外没有捕获任务活动的中间空洞。"}, - {key: "completion_tail", label: "Topology 完成尾部", unit: "milliseconds", description: "末节点 Observer 完成到 Topology 完成回调的尾部。"}, + {key: "completion_tail", label: "Topology 完成尾部", unit: "milliseconds", description: "末个真实 Taskflow 任务 completed_ms 到 Graph finished_ms 的完整尾部,内部收尾节点会进一步拆分该区间。"}, + {key: "unattributed_completion_tail", label: "未归因完成尾部", unit: "milliseconds", description: "Topology 完成尾部扣除已记录的收尾诊断区间后的剩余时间。"}, {key: "peak_entry_queue", label: "入口队列峰值", unit: "count", description: "节点进入 Worker 时看到的最大本地队列长度。"}, {key: "maximum_parallelism", label: "实际最大并行", unit: "count", description: "Observer 时间区间内同时执行的最大节点数量。"}, {key: "layer_count", label: "依赖层", unit: "count", description: "每帧 DAG 的依赖层数量。"}, @@ -1943,15 +1968,22 @@ function expanded_taskflow_graph(graph: Taskflow_Graph_Trace): Taskflow_Graph_Tr successors: successors.get(node.native_id) ?? []}))}; } -function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_state}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]; +type Taskflow_Fullscreen_View = "dag" | "timeline"; +type Taskflow_Frame_Navigation = {index: number; count: number; sequence: number; on_change: (index: number) => void}; + +function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_state, fullscreen: controlled_fullscreen, + on_fullscreen_change, on_fullscreen_view_change, frame_navigation}: {graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]; frame?: Taskflow_Frame_Trace; aggregate?: Taskflow_Graph_Aggregate | null; components: Component[]; - gallery_state?: Gallery_Pipeline_State | null}) { + gallery_state?: Gallery_Pipeline_State | null; fullscreen?: boolean; on_fullscreen_change?: (value: boolean) => void; + on_fullscreen_view_change?: (view: Taskflow_Fullscreen_View) => void; frame_navigation?: Taskflow_Frame_Navigation}) { const [nodes, set_nodes] = useState([]); const [edges, set_edges] = useState([]); const [copy_state, set_copy_state] = useState("复制拓扑 JSON"); const [viewport_width, set_viewport_width] = useState(1000); const [viewport_height, set_viewport_height] = useState(720); - const [fullscreen, set_fullscreen] = useState(false); + const [local_fullscreen, set_local_fullscreen] = useState(false); + const fullscreen = controlled_fullscreen ?? local_fullscreen; + const set_fullscreen = (value: boolean) => on_fullscreen_change ? on_fullscreen_change(value) : set_local_fullscreen(value); useEffect(() => { let cancelled = false; const displayed_graph = expanded_taskflow_graph(graph); @@ -2051,13 +2083,24 @@ function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_ } catch { set_copy_state("复制失败"); } }; - const flow = ; return
+
+ + {fullscreen && frame_navigation ?
+ + #{frame_navigation.sequence} · {frame_navigation.index + 1}/{frame_navigation.count} + +
: null} + {fullscreen && on_fullscreen_view_change ? : null} + +
{aggregate ? <>稳定节点 存在波动 @@ -2067,10 +2110,6 @@ function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_ 已执行,执行时间为主 真实排队时间大于执行时间}
-
- - -
{fullscreen ?
{flow}
: & { phase: Taskflow_Timeline_Phase; }; @@ -2099,10 +2138,17 @@ const taskflow_timeline_time_steps = { * react-calendar-timeline 使用日历毫秒。把 1 帧内毫秒放大为 1 日历秒, * 既保留微小任务的可缩放宽度,又始终以 Render_Frame 创建时刻为零点显示。 */ -function Taskflow_Timeline({graph, executions}: { - graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]; +function Taskflow_Timeline({graph, executions, frame, fullscreen: controlled_fullscreen, on_fullscreen_change, on_fullscreen_view_change, frame_navigation}: { + graph: Taskflow_Graph_Trace; executions: Taskflow_Execution_Trace[]; frame?: Taskflow_Frame_Trace; fullscreen?: boolean; + on_fullscreen_change?: (value: boolean) => void; on_fullscreen_view_change?: (view: Taskflow_Fullscreen_View) => void; + frame_navigation?: Taskflow_Frame_Navigation; }) { - const [fullscreen, set_fullscreen] = useState(false); + const [local_fullscreen, set_local_fullscreen] = useState(false); + const [copy_state, set_copy_state] = useState("复制时间线 JSON"); + const [line_height, set_line_height] = useState(52); + const [visible_ranges, set_visible_ranges] = useState>({}); + const fullscreen = controlled_fullscreen ?? local_fullscreen; + const set_fullscreen = (value: boolean) => on_fullscreen_change ? on_fullscreen_change(value) : set_local_fullscreen(value); const origin = Date.UTC(2000, 0, 1); const scale = 1000; const model = useMemo(() => { @@ -2117,68 +2163,125 @@ function Taskflow_Timeline({graph, executions}: { const items: Taskflow_Timeline_Item[] = []; const encode_time = (milliseconds: number) => origin + milliseconds * scale; const add_item = (id: string, group: string, phase: Taskflow_Timeline_Phase, - start: number, end: number, label: string, detail: string) => { + start: number, end: number, label: string, task_name: string, detail: string) => { if (!Number.isFinite(start) || !Number.isFinite(end) || end < start) return; const visible_end = Math.max(end, start + .002); + const topology_offset = start - graph.submitted_ms; items.push({id, group, phase, title: label, start_time: encode_time(start), end_time: encode_time(visible_end), canMove: false, canResize: false, canChangeGroup: false, className: `taskflowTimelineItem taskflowTimelineItem-${phase}`, - itemProps: {title: `${detail}\n+${start.toFixed(3)} → +${end.toFixed(3)} ms`}}); + itemProps: {title: `任务名称:${task_name}\n时间起点:+${start.toFixed(3)} ms(相对帧创建)\n相对 Topology:${topology_offset >= 0 ? "+" : ""}${topology_offset.toFixed(3)} ms\n时间终点:+${end.toFixed(3)} ms\n持续时间:${milliseconds(Math.max(0, end - start))}\n${detail}`}}); }; add_item("topology", "topology", "topology", graph.submitted_ms, graph.finished_ms, `Topology ${milliseconds(graph.finished_ms - graph.submitted_ms)}`, - `${graph.stage} · ${graph.name}`); + graph.name || graph.stage, `${graph.stage} · Topology 墙钟`); + const task_rows = rows.filter(row => node_by_native_id.get(row.native_id)?.type.toLowerCase() !== "diagnostic"); + const last_task_completed = task_rows.length ? Math.max(...task_rows.map(row => row.completed_ms)) : graph.finished_ms; + const last_completed = rows.length ? Math.max(...rows.map(row => row.completed_ms)) : graph.finished_ms; + if (last_completed < graph.finished_ms) { + groups.push({id: "completion-tail", title:
Topology 完成尾部 + 末任务 completed → Graph finished
, node_name: graph.name, type: "completion_tail"}); + add_item("completion-tail", "completion-tail", "completion_tail", last_completed, graph.finished_ms, + `尾部 ${milliseconds(graph.finished_ms - last_completed)}`, "Topology 完成尾部", + `${graph.stage} · 未拆分 observer/runtime graph finalization`); + } rows.forEach((execution, index) => { const node = node_by_native_id.get(execution.native_id); const group = `task-${index}`; const node_name = node?.name || execution.node_id || execution.native_id; - groups.push({id: group, node_name, worker_id: execution.worker_id, + const diagnostic = node?.type.toLowerCase() === "diagnostic"; + groups.push({id: group, node_name, worker_id: diagnostic ? undefined : execution.worker_id, started_ms: execution.started_ms, type: node?.type, title:
{taskflow_node_name(node_name)} - W{execution.worker_id} · 开始 +{execution.started_ms.toFixed(3)} ms · {node?.type ?? "task"} + {diagnostic ? "Runtime" : `W${execution.worker_id}`} · 开始 +{execution.started_ms.toFixed(3)} ms · {node?.type ?? "task"}
}); - const detail = `${node_name} · Worker ${execution.worker_id}`; - add_item(`${group}-queue`, group, "queue", execution.ready_ms, - execution.entered_ms, `排队 ${milliseconds(execution.queue_wait_ms)}`, detail); - add_item(`${group}-entry`, group, "observer_entry", execution.entered_ms, - execution.started_ms, `entry ${milliseconds(execution.observer_entry_ms)}`, detail); - add_item(`${group}-body`, group, "body", execution.started_ms, - execution.finished_ms, `执行 ${milliseconds(execution.duration_ms)}`, detail); - add_item(`${group}-exit`, group, "observer_exit", execution.finished_ms, - execution.completed_ms, `exit ${milliseconds(execution.observer_exit_ms)}`, detail); + const detail = diagnostic ? `${node_name} · Runtime tail` : `${node_name} · Worker ${execution.worker_id}`; + if (diagnostic) { + add_item(`${group}-runtime`, group, "completion_tail", execution.started_ms, + execution.finished_ms, `${taskflow_node_name(node_name)} ${milliseconds(execution.duration_ms)}`, + node_name, detail); + return; + } + add_item(`${group}-queue`, group, "queue", execution.ready_ms, execution.entered_ms, + `排队 ${milliseconds(execution.queue_wait_ms)}`, node_name, detail); + add_item(`${group}-entry`, group, "observer_entry", execution.entered_ms, execution.started_ms, + `entry ${milliseconds(execution.observer_entry_ms)}`, node_name, detail); + add_item(`${group}-body`, group, "body", execution.started_ms, execution.finished_ms, + `执行 ${milliseconds(execution.duration_ms)}`, node_name, detail); + add_item(`${group}-exit`, group, "observer_exit", execution.finished_ms, execution.completed_ms, + `exit ${milliseconds(execution.observer_exit_ms)}`, node_name, detail); }); const extent = Math.max(.1, graph.finished_ms, ...rows.flatMap(row => [row.completed_ms, row.finished_ms])); return {groups, items, end: origin + extent * 1.06 * scale, - span: Math.max(.1, extent * 1.06) * scale, executions: rows.length}; + span: Math.max(.1, extent * 1.06) * scale, executions: rows.length, last_completed, last_task_completed}; }, [graph, executions]); + const timeline_view_key = `${graph.stage}:${graph.name}`; + const visible_range = visible_ranges[timeline_view_key]; + const timeline_time_props = visible_range + ? {visibleTimeStart: visible_range.start, visibleTimeEnd: visible_range.end} + : {defaultTimeStart: origin, defaultTimeEnd: model.end}; + const copy_timeline = async () => { + const native_ids = new Set(graph.nodes.map(node => node.native_id)); + try { + await copy_text(JSON.stringify({ + sequence: frame?.sequence, + correlation_id: frame?.correlation_id, + markers: frame?.markers ?? {}, + graph, + timeline: {origin: "render_frame_created", submitted_ms: graph.submitted_ms, finished_ms: graph.finished_ms, + wall_time_ms: Math.max(0, graph.finished_ms - graph.submitted_ms), last_execution_completed_ms: model.last_completed, + last_task_completed_ms: model.last_task_completed, + completion_tail_ms: Math.max(0, graph.finished_ms - model.last_task_completed)}, + 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
+
+ + {fullscreen && frame_navigation ?
+ + #{frame_navigation.sequence} · {frame_navigation.index + 1}/{frame_navigation.count} + +
: null} + {fullscreen && on_fullscreen_view_change ? : null} + + + + +
Topology 墙钟 Executor 排队 Observer 任务体执行 + Topology 完成尾部 零点为 Render_Frame 创建时刻 · {model.executions} 次执行
-
- key={`${graph.stage}:${graph.submitted_ms}:${fullscreen ? "fullscreen" : "normal"}`} - groups={model.groups} items={model.items} keys={taskflow_timeline_keys} - defaultTimeStart={origin} defaultTimeEnd={model.end} + key={timeline_view_key} groups={model.groups} items={model.items} keys={taskflow_timeline_keys} + {...timeline_time_props} + onTimeChange={(start, end, update_scroll_canvas) => { + set_visible_ranges(current => ({...current, [timeline_view_key]: {start, end}})); + update_scroll_canvas(start, end); + }} sidebarWidth={fullscreen ? 420 : 340} rightSidebarWidth={0} - lineHeight={52} itemHeightRatio={.68} itemVerticalGap={5} + lineHeight={line_height} itemHeightRatio={.68} itemVerticalGap={5} minZoom={10} maxZoom={Math.max(model.span * 20, 1000)} buffer={1} canMove={false} canResize={false} canChangeGroup={false} canSelect={false} stackItems={false} traditionalZoom timeSteps={taskflow_timeline_time_steps} groupRenderer={({group}) => group.title}> - + {({getRootProps}) =>
任务 / Worker / 开始偏移
}
"相对帧创建时刻的偏移"}/> @@ -2206,6 +2309,7 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon const [frame_index, set_frame_index] = useState(0); const [stage_key, set_stage_key] = useState(""); const [view_mode, set_view_mode] = useState<"aggregate" | "single" | "timeline">("aggregate"); + const [fullscreen_view, set_fullscreen_view] = useState(null); const [busy, set_busy] = useState(false); const [error, set_error] = useState(""); const [gallery_state, set_gallery_state] = useState(null); @@ -2275,6 +2379,20 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon const graph_analysis = useMemo(() => graph && frame ? taskflow_graph_analysis(graph, frame.executions) : null, [graph, frame]); const paint_analysis = useMemo(() => frame_paint_analysis(frame), [frame]); + const displayed_view_mode = fullscreen_view === "dag" ? "single" : fullscreen_view === "timeline" ? "timeline" : view_mode; + const switch_fullscreen_view = (next: Taskflow_Fullscreen_View) => { + set_view_mode(next === "dag" ? "single" : "timeline"); + set_fullscreen_view(next); + }; + const change_frame = (next_index: number) => { + if (!response?.frames.length) return; + const index = Math.max(0, Math.min(response.frames.length - 1, next_index)); + const next_frame = response.frames[index]; + set_frame_index(index); + if (!next_frame.graphs.some(value => taskflow_stage_key(value) === stage_key)) + set_stage_key(next_frame.graphs[0] ? taskflow_stage_key(next_frame.graphs[0]) : ""); + }; + const frame_navigation = response && frame ? {index: frame_index, count: response.frames.length, sequence: frame.sequence, on_change: change_frame} : undefined; return
void load()}/>
{error ?

{error}

: null} {frame ? <> -
@@ -2296,7 +2414,7 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon
Executor {frame.worker_count} workers · 本帧 {frame.executions.length} 次任务执行
- {view_mode === "aggregate" && aggregate ? <>
+ {displayed_view_mode === "aggregate" && aggregate ? <>
业务阶段
{aggregate.graph.stage}
捕获帧 / 阶段样本
{aggregate.frames} / {aggregate.samples}
拓扑变体
{aggregate.topology_variants}
@@ -2317,9 +2435,10 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon
{definition.label}
{taskflow_distribution_value(statistic, "milliseconds")}
; }) : null} - : view_mode === "single" && graph && graph_analysis ? <>
-
业务阶段
{graph.stage}
DAG 节点
{graph.nodes.length}
-
Topology 墙钟
{milliseconds(graph_analysis.wall_time)}
+ : displayed_view_mode === "single" && graph && graph_analysis ? <>
+
业务阶段
{graph.stage}
Taskflow 节点
{graph_analysis.node_count}
+
收尾诊断节点
{graph_analysis.diagnostic_node_count}
+
Topology 墙钟
{milliseconds(graph_analysis.wall_time)}
叶子任务墙钟总和
{milliseconds(graph_analysis.execution_time)}
Module 包络总和
{milliseconds(graph_analysis.module_envelope_time)}
叶子任务线程 CPU{graph_analysis.cpu_time_coarse ? "(低分辨率)" : ""}
{milliseconds(graph_analysis.cpu_execution_time)}
@@ -2335,7 +2454,13 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon
累计节点排队
{milliseconds(graph_analysis.queue_time)}
首次准入等待
{milliseconds(graph_analysis.initial_wait)}
内部调度空洞
{milliseconds(graph_analysis.internal_idle_time)}
-
Topology 完成尾部
{milliseconds(graph_analysis.completion_tail)}
+
Topology 完成尾部
{milliseconds(graph_analysis.completion_tail)}
+
未归因完成尾部
{milliseconds(graph_analysis.unattributed_completion_tail)}
+ {graph_analysis.diagnostic_rows.map(row => { + const node = graph.nodes.find(value => value.native_id === row.native_id); + return
+
{taskflow_node_name(node?.name ?? row.node_id)}
{milliseconds(row.duration_ms)}
; + })}
入口队列峰值
{graph_analysis.peak_entry_queue}
实际最大并行
{graph_analysis.maximum_parallelism}
依赖层 / 并行层
{graph_analysis.layer_count} / {graph_analysis.parallel_layer_count}
@@ -2350,9 +2475,13 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon
Paint 阶段衔接
{milliseconds(paint_analysis.coordination)}
: null}
状态
{graph.completed ? "完成" : "未完成"}
-
- : view_mode === "timeline" && graph - ? +
set_fullscreen_view(value ? "dag" : null)} + on_fullscreen_view_change={switch_fullscreen_view} frame_navigation={frame_navigation}/> + : displayed_view_mode === "timeline" && graph + ? set_fullscreen_view(value ? "timeline" : null)} + on_fullscreen_view_change={switch_fullscreen_view} frame_navigation={frame_navigation}/> :
该帧没有所选 Taskflow 阶段切换到“多帧聚合拓扑”仍可查看其他帧中的同一业务阶段。
} :
等待逐帧 Taskflow 样本输入 N 后捕获后续实际渲染帧;每帧独立绑定其 DAG 元信息和 Observer 结果。
}
; diff --git a/webapp_gallery/src/styles.css b/webapp_gallery/src/styles.css index 7741c02..10d0abc 100644 --- a/webapp_gallery/src/styles.css +++ b/webapp_gallery/src/styles.css @@ -202,10 +202,13 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .taskflowGraphSummary dt, .taskflowRuntimeSummary dt { color: #71839e; font-size: 10px; } .taskflowGraphSummary dd, .taskflowRuntimeSummary dd { margin: 5px 0 0; color: #5ce4c2; font: 700 13px/1.35 ui-monospace, monospace; white-space: normal; overflow-wrap: anywhere; } .taskflowDagSection { min-width: 640px; overflow: visible; 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 { display: flex; align-items: center; justify-content: flex-start; flex-wrap: wrap; 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; } -.taskflowDagActions { display: flex; flex: none; align-items: center; gap: 9px; } -.taskflowLegend { display: flex; align-items: center; flex-wrap: wrap; gap: 12px; color: #8298b4; font-size: 10px; } +.taskflowDagToolbar button:disabled { opacity: .42; cursor: default; } +.taskflowDagActions { display: flex; flex: none; align-items: center; flex-wrap: wrap; gap: 9px; } +.taskflowFullscreenFrameNav { display: inline-flex; align-items: center; gap: 7px; padding: 2px 5px; border: 1px solid #29435e; border-radius: 7px; background: #091321; } +.taskflowFullscreenFrameNav span { min-width: 112px; color: #91a5c0; font: 10px/1 ui-monospace, monospace; text-align: center; } +.taskflowLegend { display: flex; flex: 1 1 520px; 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; } @@ -278,6 +281,7 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .taskflowTimelineSidebarHeader { display: flex; align-items: center; height: 100%; padding: 0 10px; color: #91a5c0; background: #101d2f; font: 10px/1.2 ui-monospace, monospace; } .taskflowTimelineViewport .rct-sidebar { border-color: #29405e; background: #091321; } .taskflowTimelineViewport .rct-sidebar .rct-sidebar-row { overflow: visible; border-color: #1d304a; } +.taskflowTimelineHeaders { position: sticky !important; top: 0; z-index: 100; background: #0c1727; } .taskflowTimelineViewport .rct-calendar-header { border-color: #29405e; background: #0c1727; } .taskflowTimelineViewport .rct-dateHeader { color: #91a5c0; border-color: #29405e; background: #101d2f; font: 9px/1 ui-monospace, monospace; } .taskflowTimelineViewport .rct-horizontal-lines .rct-hl-even, @@ -289,6 +293,7 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .taskflowTimelineItem-queue { background: #8b5c24 !important; } .taskflowTimelineItem-observer_entry, .taskflowTimelineItem-observer_exit { background: #67468b !important; } .taskflowTimelineItem-body { background: #197462 !important; } +.taskflowTimelineItem-completion_tail { background: #78354b !important; } .taskflowTimelineMarker { width: 2px !important; z-index: 8; pointer-events: none; } .taskflowTimelineMarkerSubmitted { background: #e0a65d; } .taskflowTimelineMarkerFinished { background: #ef6a82; } @@ -296,6 +301,7 @@ canvas { display: block; width: 100%; height: 100%; background: #070d18; } .taskflowTimelineLegendQueue { border-color: #c0843e !important; background: #8b5c24 !important; } .taskflowTimelineLegendObserver { border-color: #9d74c8 !important; background: #67468b !important; } .taskflowTimelineLegendBody { border-color: #35ad92 !important; background: #197462 !important; } +.taskflowTimelineLegendCompletion { border-color: #b75d79 !important; background: #78354b !important; } .taskflowTimelineFullscreen { position: fixed; inset: 0; z-index: 10000; display: flex; flex-direction: column; min-width: 0; overflow: hidden; border: 0; border-radius: 0; background: #07101c; } .taskflowTimelineFullscreen .taskflowTimelineViewport { flex: 1; max-height: none; min-height: 0; }