From 60452211741edb4a47f9222f48d0b0cd9a03304e Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Sat, 29 Aug 2026 12:03:22 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=B9=E6=96=B9=E5=BC=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kernel/src/kernel/render_common.cpp | 79 ++++++-- kernel/src/kernel/render_common.hpp | 7 +- kernel/src/test/render_test.cpp | 10 + mcp/core/Control_Service.cpp | 11 +- mcp/core/runtime/Plot.cpp | 123 +----------- mcp/core/runtime/Plot.hpp | 5 - mcp/tests/Control_Path_Benchmarks.cpp | 12 +- web_server/src/Gallery_Video_Stream.cpp | 189 ++++++++++++++++-- web_server/src/Gallery_Video_Stream.hpp | 2 + web_server/src/Gallery_Video_Stream.ipp | 17 +- web_server/src/Gallery_WebSocket.cpp | 11 + web_server/src/WebRtc_Video_Session.cpp | 118 ++++++----- web_server/src/WebRtc_Video_Session.hpp | 8 +- web_server/src/Web_Server.cpp | 41 ++-- .../src/frame_sampling/Frame_Sampler.cpp | 23 ++- .../src/frame_sampling/Frame_Sampler.hpp | 1 + web_server/tests/Frame_Sampler_Tests.cpp | 4 +- webapp_gallery/src/app.tsx | 98 ++++++--- 18 files changed, 478 insertions(+), 281 deletions(-) diff --git a/kernel/src/kernel/render_common.cpp b/kernel/src/kernel/render_common.cpp index 9a5b18f..793f9e0 100644 --- a/kernel/src/kernel/render_common.cpp +++ b/kernel/src/kernel/render_common.cpp @@ -89,6 +89,7 @@ private: std::atomic_size_t named_task_count{}; std::atomic_size_t active_depth{}; std::atomic_size_t peak_active_depth{}; + std::atomic_bool executing_task_segment{}; std::atomic_size_t current_queue_size{}; std::atomic_size_t current_queue_capacity{}; std::atomic_size_t peak_queue_size{}; @@ -106,6 +107,8 @@ private: std::atomic_bool active_task_reported{}; std::atomic_uint64_t task_time_ns{}; std::atomic_uint64_t busy_time_ns{}; + std::atomic_uint64_t cooperative_wait_count{}; + std::atomic_uint64_t cooperative_wait_time_ns{}; std::atomic_uint64_t min_task_time_ns{std::numeric_limits::max()}; std::atomic_uint64_t max_task_time_ns{}; std::atomic_uint64_t first_task_time_ns{}; @@ -123,7 +126,8 @@ private: Clock::time_point cooperative_wait_started{}; /* 主动 corun 让出 Worker 的墙钟起点。 */ std::uint64_t cooperative_wait_ns{}; /* 已累计 cooperative wait 墙钟。 */ std::array cooperative_waits{}; /* 按帧时间线的固定容量让出区间。 */ - std::size_t cooperative_wait_count{}; + std::uint64_t cooperative_wait_count{}; + std::size_t traced_cooperative_wait_count{}; Render_Frame* frame{}; /* 进入任务时唯一活动的按帧捕获。 */ std::size_t queue_size{}; /* 进入任务时 worker 队列深度。 */ std::size_t queue_capacity{}; /* 进入任务时 worker 队列容量。 */ @@ -134,6 +138,10 @@ private: std::vector> starts; std::unique_ptr worker_statistics; std::size_t worker_statistics_count{}; + std::atomic_size_t active_task_count{}; + std::atomic_size_t peak_active_task_count{}; + std::atomic_size_t active_worker_count{}; + std::atomic_size_t peak_active_worker_count{}; std::uint64_t worker_occupation_limit_ns{}; /* 单节点连续非 CPU 等待 Worker 的上限。 */ Task_Overrun_Action worker_overrun_action{}; /* 节点超过占用上限后的处置策略。 */ std::atomic_bool watchdog_stopping{}; /* Watchdog 生命周期停止标志。 */ @@ -156,6 +164,18 @@ private: auto current = value.load(std::memory_order_relaxed); while ((current == 0 || next < current) && !value.compare_exchange_weak(current, next, std::memory_order_relaxed)) {} } + void set_worker_execution(Worker_Statistics& worker, bool executing) noexcept { + const auto previous = worker.executing_task_segment.exchange( + executing, std::memory_order_acq_rel); + if (previous == executing) return; + if (!executing) { + active_worker_count.fetch_sub(1, std::memory_order_relaxed); + return; + } + const auto active = active_worker_count.fetch_add( + 1, std::memory_order_relaxed) + 1; + update_max(peak_active_worker_count, active); + } static std::size_t task_type_index(tf::TaskType type) noexcept { return static_cast(type); } @@ -257,13 +277,15 @@ private: active.cooperative_wait_started == Clock::time_point{}; if (begins_wait) active.cooperative_wait_started = now; + if (begins_wait) ++active.cooperative_wait_count; if (begins_wait && active.frame != nullptr && - active.cooperative_wait_count < active.cooperative_waits.size()) { - active.cooperative_waits[active.cooperative_wait_count++] = {now, {}}; + active.traced_cooperative_wait_count < active.cooperative_waits.size()) { + active.cooperative_waits[active.traced_cooperative_wait_count++] = {now, {}}; } active.cooperatively_suspended = true; worker_statistics[worker].active_segment_started_ns.store( 0, std::memory_order_release); + set_worker_execution(worker_statistics[worker], false); #if defined(_WIN32) worker_statistics[worker].active_cpu_cycles.store( 0, std::memory_order_release); @@ -280,9 +302,9 @@ private: std::chrono::duration_cast( now - active.cooperative_wait_started).count()); active.cooperative_wait_started = {}; - if (active.cooperative_wait_count != 0) { + if (active.traced_cooperative_wait_count != 0) { auto& wait = active.cooperative_waits[ - active.cooperative_wait_count - 1]; + active.traced_cooperative_wait_count - 1]; if (wait.second == Clock::time_point{}) wait.second = now; } } @@ -290,6 +312,7 @@ private: active.segment_started = now; worker_statistics[worker].active_segment_started_ns.store( clock_ns(now), std::memory_order_release); + set_worker_execution(worker_statistics[worker], true); #if defined(_WIN32) worker_statistics[worker].active_cpu_cycles.store(0, std::memory_order_release); @@ -369,9 +392,13 @@ public: worker_starts.push_back(std::move(record)); const auto 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; + const auto previous_depth = worker_state.active_depth.fetch_add( + 1, std::memory_order_relaxed); + const auto active_depth = previous_depth + 1; update_max(worker_state.peak_active_depth, active_depth); + const auto active_tasks = active_task_count.fetch_add( + 1, std::memory_order_relaxed) + 1; + update_max(peak_active_task_count, active_tasks); worker_state.current_queue_size.store(queue_size, std::memory_order_relaxed); worker_state.current_queue_capacity.store(queue_capacity, std::memory_order_relaxed); worker_state.active_task_hash.store(task.hash_value(), std::memory_order_relaxed); @@ -391,6 +418,7 @@ public: update_first(worker_state.first_task_time_ns, clock_ns(now)); worker_starts.back().started = Clock::now(); worker_starts.back().segment_started = worker_starts.back().started; + set_worker_execution(worker_state, true); } void on_exit(tf::WorkerView worker, tf::TaskView task) override { const auto finished = Clock::now(); @@ -401,9 +429,9 @@ public: std::chrono::duration_cast( finished - active.cooperative_wait_started).count()); active.cooperative_wait_started = {}; - if (active.cooperative_wait_count != 0) { + if (active.traced_cooperative_wait_count != 0) { auto& wait = active.cooperative_waits[ - active.cooperative_wait_count - 1]; + active.traced_cooperative_wait_count - 1]; if (wait.second == Clock::time_point{}) wait.second = finished; } } @@ -428,6 +456,10 @@ public: worker_state.task_time_ns.fetch_add(occupied, std::memory_order_relaxed); worker_state.busy_time_ns.fetch_add( start.active_time_ns, std::memory_order_relaxed); + worker_state.cooperative_wait_count.fetch_add( + start.cooperative_wait_count, std::memory_order_relaxed); + worker_state.cooperative_wait_time_ns.fetch_add( + start.cooperative_wait_ns, std::memory_order_relaxed); update_min(worker_state.min_task_time_ns, occupied); update_max(worker_state.max_task_time_ns, occupied); auto type_index = task_type_index(task.type()); @@ -446,6 +478,7 @@ public: worker_state.longest_task_type.store(task.type(), std::memory_order_relaxed); } worker_state.active_depth.fetch_sub(1, std::memory_order_relaxed); + active_task_count.fetch_sub(1, std::memory_order_relaxed); std::optional trace_task; if (start.frame) { try { @@ -455,7 +488,7 @@ public: start.queue_size, start.queue_capacity, start.entered, start.started, finished, start.cooperative_wait_ns, std::span{start.cooperative_waits}.first( - start.cooperative_wait_count)); + start.traced_cooperative_wait_count)); } catch (...) { /* Observer 不能让按需诊断分配失败改变渲染任务的完成语义。 */ @@ -468,6 +501,7 @@ public: *start.frame, worker.id(), *trace_task, completed); } if (worker_starts.empty()) { + set_worker_execution(worker_state, false); worker_state.active_task_hash.store(0, std::memory_order_relaxed); worker_state.active_task_started_ns.store(0, std::memory_order_relaxed); worker_state.active_segment_started_ns.store(0, @@ -488,6 +522,7 @@ public: worker_state.active_task_started_ns.store(clock_ns(parent.entered), std::memory_order_relaxed); if (parent.cooperatively_suspended) { + set_worker_execution(worker_state, false); worker_state.active_segment_started_ns.store( 0, std::memory_order_release); #if defined(_WIN32) @@ -498,6 +533,7 @@ public: #endif } else { + set_worker_execution(worker_state, true); const auto parent_resumed = Clock::now(); parent.segment_started = parent_resumed; worker_state.active_segment_started_ns.store( @@ -540,6 +576,8 @@ public: state.total_task_time_ns = 0; state.worker_busy_time_ns = 0; state.worker_cpu_time_ns = 0; + state.cooperative_wait_count = 0; + state.cooperative_wait_time_ns = 0; state.longest_task_time_ns = 0; state.longest_task_hash = 0; state.longest_task_name.clear(); @@ -574,6 +612,10 @@ public: : std::string{}; target.task_time_ns = source.task_time_ns.load(std::memory_order_relaxed); target.busy_time_ns = source.busy_time_ns.load(std::memory_order_relaxed); + target.cooperative_wait_count = source.cooperative_wait_count.load( + std::memory_order_relaxed); + target.cooperative_wait_time_ns = source.cooperative_wait_time_ns.load( + std::memory_order_relaxed); #if defined(_WIN32) const auto native_handle = source.native_thread_handle.load( std::memory_order_acquire); @@ -585,9 +627,6 @@ public: #else target.cpu_time_ns = 0; #endif - target.non_cpu_time_ns = target.busy_time_ns > target.cpu_time_ns - ? target.busy_time_ns - target.cpu_time_ns - : 0; target.idle_time_ns = state.observed_wall_time_ns > target.busy_time_ns ? state.observed_wall_time_ns - target.busy_time_ns : 0; auto min = source.min_task_time_ns.load(std::memory_order_relaxed); target.min_task_time_ns = target.task_count ? min : 0; @@ -597,12 +636,6 @@ public: ? static_cast(target.cpu_time_ns) * 100.0 / static_cast(state.observed_wall_time_ns) : 0.0; - const auto active_depth = source.active_depth.load( - std::memory_order_relaxed); - state.active_task_count += active_depth; - state.peak_active_task_count += source.peak_active_depth.load( - std::memory_order_relaxed); - state.active_worker_count += active_segment != 0; state.observed_task_count += target.task_count; state.named_task_count += source.named_task_count.load( std::memory_order_relaxed); @@ -615,6 +648,8 @@ public: state.total_task_time_ns += target.task_time_ns; state.worker_busy_time_ns += target.busy_time_ns; state.worker_cpu_time_ns += target.cpu_time_ns; + state.cooperative_wait_count += target.cooperative_wait_count; + state.cooperative_wait_time_ns += target.cooperative_wait_time_ns; const auto worker_first = source.first_task_time_ns.load( std::memory_order_relaxed); if (worker_first && (!first || worker_first < first)) first = worker_first; @@ -651,8 +686,10 @@ public: } } state.observed_wall_time_ns = first && last >= first ? last - first : 0; - state.peak_active_worker_count = std::max( - state.peak_active_worker_count, state.active_worker_count); + state.active_task_count = active_task_count.load(std::memory_order_relaxed); + state.peak_active_task_count = peak_active_task_count.load(std::memory_order_relaxed); + state.active_worker_count = active_worker_count.load(std::memory_order_relaxed); + state.peak_active_worker_count = peak_active_worker_count.load(std::memory_order_relaxed); state.worker_utilization = workers && state.observed_wall_time_ns ? static_cast(state.worker_busy_time_ns) * 100.0 / static_cast(state.observed_wall_time_ns) / diff --git a/kernel/src/kernel/render_common.hpp b/kernel/src/kernel/render_common.hpp index 7657df2..291b485 100644 --- a/kernel/src/kernel/render_common.hpp +++ b/kernel/src/kernel/render_common.hpp @@ -52,8 +52,9 @@ struct Task_Worker_State { std::string active_task_type{}; /* 当前任务的 Taskflow 原生 TaskType;空闲时为空。 */ std::uint64_t task_time_ns{}; /* 已完成任务体墙钟之和,不含 cooperative wait。 */ std::uint64_t busy_time_ns{}; /* Worker 实际执行任务体片段的墙钟并集,不含 cooperative wait。 */ - std::uint64_t cpu_time_ns{}; /* Worker最外层任务活跃区间内累计的线程 CPU 时间。 */ - std::uint64_t non_cpu_time_ns{}; /* busy_time_ns 减去 cpu_time_ns;只表示未计费墙钟,不推断锁或抢占。 */ + std::uint64_t cpu_time_ns{}; /* 从首次观察任务起,Worker 原生线程的累计 CPU 时间。 */ + std::uint64_t cooperative_wait_count{}; /* 当前 Worker 上任务主动 corun/corun_until 让出的累计次数。 */ + std::uint64_t cooperative_wait_time_ns{}; /* 外层任务协作挂起区间之和;期间 Worker 可执行其他任务。 */ std::uint64_t idle_time_ns{}; std::uint64_t min_task_time_ns{}; std::uint64_t max_task_time_ns{}; @@ -87,6 +88,8 @@ struct Task_Runtime_State : State_Type { std::uint64_t total_task_time_ns{}; std::uint64_t worker_busy_time_ns{}; std::uint64_t worker_cpu_time_ns{}; + std::uint64_t cooperative_wait_count{}; + std::uint64_t cooperative_wait_time_ns{}; std::uint64_t observed_wall_time_ns{}; double worker_utilization{}; double worker_cpu_utilization{}; diff --git a/kernel/src/test/render_test.cpp b/kernel/src/test/render_test.cpp index af8becd..73844fa 100644 --- a/kernel/src/test/render_test.cpp +++ b/kernel/src/test/render_test.cpp @@ -2,6 +2,7 @@ #include "Frame_Policy/Frame_Policy.hpp" #include #include +#include #include namespace { struct Direct_Renderable : double_buffer::Def { @@ -271,6 +272,15 @@ TEST(task_graph_observer, business_dag_and_native_execution_share_node_identity) EXPECT_GE(cooperative_execution->cooperative_wait_ms, 0.0); EXPECT_GE(cooperative_execution->cooperative_waits.front().finished_ms, cooperative_execution->cooperative_waits.front().started_ms); + const auto runtime = aethera::task_runtime_state(); + EXPECT_GE(runtime.cooperative_wait_count, 1U); + EXPECT_LE(runtime.active_worker_count, runtime.worker_count); + EXPECT_EQ(runtime.cooperative_wait_count, + std::accumulate(runtime.workers.begin(), runtime.workers.end(), + std::uint64_t{}, [](std::uint64_t total, + const auto& worker) { + return total + worker.cooperative_wait_count; + })); std::unordered_set metadata_ids; for (const auto& node : trace.graphs.front().nodes) diff --git a/mcp/core/Control_Service.cpp b/mcp/core/Control_Service.cpp index e90a6f2..15d7a1d 100644 --- a/mcp/core/Control_Service.cpp +++ b/mcp/core/Control_Service.cpp @@ -201,8 +201,8 @@ template for (const auto& worker : state.workers) workers.push_back({ {"id", worker.id}, {"task_count", worker.task_count}, - {"current_queue_size", worker.current_queue_size}, - {"current_queue_capacity", worker.current_queue_capacity}, + {"entry_queue_size", worker.current_queue_size}, + {"entry_queue_capacity", worker.current_queue_capacity}, {"peak_queue_size", worker.peak_observed_queue_size}, {"max_queue_capacity", worker.max_observed_queue_capacity}, {"active_task", {{"native_id", std::to_string(worker.active_task_hash)}, @@ -211,7 +211,8 @@ template {"task_time_ns", worker.task_time_ns}, {"busy_time_ns", worker.busy_time_ns}, {"cpu_time_ns", worker.cpu_time_ns}, - {"non_cpu_time_ns", worker.non_cpu_time_ns}, + {"cooperative_wait_count", worker.cooperative_wait_count}, + {"cooperative_wait_time_ns", worker.cooperative_wait_time_ns}, {"idle_time_ns", worker.idle_time_ns}, {"min_task_time_ns", worker.min_task_time_ns}, {"max_task_time_ns", worker.max_task_time_ns}, @@ -225,7 +226,7 @@ template {"min_time_ns", type.min_time_ns}, {"max_time_ns", type.max_time_ns}}); return { - {"protocol", "aethera.taskflow.runtime"}, {"version", 1}, + {"protocol", "aethera.taskflow.runtime"}, {"version", 2}, {"worker_count", state.worker_count}, {"active_topologies", state.active_topology_count}, {"active_taskflows", state.active_taskflow_count}, @@ -247,6 +248,8 @@ template {"total_task_time_ns", state.total_task_time_ns}, {"worker_busy_time_ns", state.worker_busy_time_ns}, {"worker_cpu_time_ns", state.worker_cpu_time_ns}, + {"cooperative_wait_count", state.cooperative_wait_count}, + {"cooperative_wait_time_ns", state.cooperative_wait_time_ns}, {"observed_wall_time_ns", state.observed_wall_time_ns}, {"worker_utilization", state.worker_utilization}, {"worker_cpu_utilization", state.worker_cpu_utilization}, diff --git a/mcp/core/runtime/Plot.cpp b/mcp/core/runtime/Plot.cpp index b53b5c4..9a919ff 100644 --- a/mcp/core/runtime/Plot.cpp +++ b/mcp/core/runtime/Plot.cpp @@ -629,15 +629,6 @@ struct Plot::Private { std::atomic_uint64_t taskflow_trace_control{}; std::array>, maximum_taskflow_trace_frames> taskflow_trace_slots{}; - std::atomic_size_t post_publish_trace_remaining{}; /* 仅捕获 publish 后外接 DAG 的剩余样本。 */ - std::atomic_uint64_t post_publish_trace_control{}; /* 高 32 位 requested,低 32 位 captured。 */ - std::array>, - maximum_taskflow_trace_frames> post_publish_trace_slots{}; - Task_Graph post_publish_graph{"plot.post_publish"}; /* publish 后外接 DAG;不再占用 Scene render admission。 */ - Task_Node post_publish_tail{}; - bool has_post_publish_tail{}; - std::atomic_bool post_publish_busy{}; - std::vector> completion_extensions{}; /* 生命周期覆盖 post-publish module 借用。 */ Frame_Statistics_Accumulator completed_frame_statistics{diagnostic_window_capacity}; std::uint64_t applied_statistics_generation{}; /* 仅完成帧退役任务读写。 */ std::atomic_uint64_t statistics_generation{}; /* reset 只推进代次,不触碰单写者累加器。 */ @@ -700,9 +691,7 @@ struct Plot::Private { void retire_completed_frame(not_null frame); void consume_retired_frames(); void finalize_retired_frame(not_null frame); - void attach_completion(std::unique_ptr completion); [[nodiscard]] bool mark_taskflow_trace(Render_Frame& frame); - [[nodiscard]] bool mark_post_publish_taskflow_trace(); void store_trace(std::atomic_uint64_t& control, std::array>, maximum_taskflow_trace_frames>& slots, @@ -962,16 +951,6 @@ bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) { } return false; } -bool Plot::Private::mark_post_publish_taskflow_trace() { - auto remaining = post_publish_trace_remaining.load(std::memory_order_acquire); - while (remaining != 0) { - if (post_publish_trace_remaining.compare_exchange_weak( - remaining, remaining - 1, std::memory_order_acq_rel, - std::memory_order_acquire)) - return true; - } - return false; -} void Plot::Private::store_trace( std::atomic_uint64_t& control, std::array>, @@ -1335,61 +1314,10 @@ void Plot::Private::consume_completed_frame(not_null frame) { break; } if (!managed) throw std::logic_error("frame callback has no owned Plot frame"); - /* - * Scene 已在调用本 callback 前释放自己的 render admission;这里同步 - * 释放 Plot 的 view/update 门,并立刻唤醒 busy 期间保留的 latest tick。 - * 之后外接 DAG 仍在当前 Frame trace 内执行,但不会阻塞下一帧渲染。 - */ - if (std::holds_alternative>(managed->frame)) release_render_admission(lifetime); - if (post_publish_graph.empty()) return; - bool expected = false; - if (!post_publish_busy.compare_exchange_strong( - expected, true, std::memory_order_acq_rel, - std::memory_order_acquire)) - return; - /* - * 外接图异步提交;上一轮尚未完成时直接合并到 sampler 内的 latest,绝不 - * 在 Taskflow Worker 内等待。诊断使用独立 Render_Frame 保存外接图观察 - * 窗口,因此 Scene 的物理帧可立即退役并被下一次渲染复用。 - */ - const bool local_post_publish_trace = mark_post_publish_taskflow_trace(); - auto trace_frame = local_post_publish_trace - ? std::make_shared(frame->identity()) - : std::shared_ptr{}; - bool local_trace_started{}; - if (trace_frame) { - trace_frame->request_taskflow_trace(); - local_trace_started = - aethera::detail::begin_taskflow_trace(*trace_frame); - } - const auto weak = lifetime; - auto completion = [weak, trace_frame, local_trace_started] { - const auto owner = weak.lock(); - if (!owner) return; - if (local_trace_started) { - aethera::detail::finish_taskflow_trace(*trace_frame); - owner->d->store_trace( - owner->d->post_publish_trace_control, - owner->d->post_publish_trace_slots, - trace_frame->take_taskflow_trace()); - } - owner->d->post_publish_busy.store(false, std::memory_order_release); - }; - try { - if (trace_frame) - aethera::detail::run_taskflow( - post_publish_graph, *trace_frame, "plot.post_publish", - std::move(completion)); - else - aethera::detail::run_taskflow( - post_publish_graph, std::move(completion)); - } - catch (...) { - if (local_trace_started) aethera::detail::finish_taskflow_trace(*trace_frame); - if (local_post_publish_trace) post_publish_trace_remaining.fetch_add(1, std::memory_order_release); - post_publish_busy.store(false, std::memory_order_release); - throw; - } + /* Scene 完成即释放 Plot admission。媒体采样属于 Gallery 自己的独立时钟和 + * DAG,不再借 Plot 保存一套 post-publish 管线状态。 */ + if (std::holds_alternative>(managed->frame)) + release_render_admission(lifetime); } void Plot::Private::consume_retired_frames() { for (;;) { @@ -1533,27 +1461,11 @@ void Plot::Private::finalize_retired_frame(not_null frame) { }); } } -void Plot::Private::attach_completion( - std::unique_ptr completion) { - if (!completion || completion->empty()) throw std::invalid_argument("Plot completion pipeline is empty"); - completion_extensions.push_back(std::move(completion)); - auto extension = post_publish_graph.compose( - completion_extensions.back()->name(), *completion_extensions.back()); - extension.describe("owner", "plot") - .describe("stage", "post-publish frame pipeline extension"); - if (has_post_publish_tail) post_publish_tail.precede(extension); - post_publish_tail = std::move(extension); - has_post_publish_tail = true; -} Plot::Plot(std::unique_ptr scene, std::unique_ptr view) : d(std::make_unique(std::move(scene), std::move(view))) {} Plot::Plot(std::unique_ptr scene, std::unique_ptr view) : d(std::make_unique(std::move(scene), std::move(view))) {} Plot::~Plot() = default; -void Plot::attach_scene_completion( - std::unique_ptr completion) { - d->attach_completion(std::move(completion)); -} void Plot::ensure_started() { std::call_once(d->start_once, [this] { const auto weak = weak_from_this(); @@ -1900,33 +1812,6 @@ nlohmann::json Plot::taskflow_trace() const { d->taskflow_trace_remaining, d->taskflow_trace_slots); } -void Plot::request_post_publish_taskflow_trace(std::size_t frame_count) { - if (frame_count == 0 || - frame_count > Private::maximum_taskflow_trace_frames) - throw std::invalid_argument( - "Taskflow post-publish trace frame_count must be between 1 and 120"); - ensure_started(); - auto control = d->post_publish_trace_control.load(std::memory_order_acquire); - for (;;) { - const auto requested = static_cast(control >> 32U); - const auto captured = static_cast(control); - if (requested != captured) - throw std::logic_error( - "A post-publish Taskflow trace request is already active"); - const auto next = static_cast(frame_count) << 32U; - if (d->post_publish_trace_control.compare_exchange_weak( - control, next, std::memory_order_release, - std::memory_order_acquire)) - break; - } - for (auto& slot : d->post_publish_trace_slots) slot.store({}, std::memory_order_release); - d->post_publish_trace_remaining.store(frame_count, std::memory_order_release); -} -nlohmann::json Plot::post_publish_taskflow_trace() const { - return d->trace_response(d->post_publish_trace_control, - d->post_publish_trace_remaining, - d->post_publish_trace_slots); -} void Plot::reset_diagnostics() { std::visit([](auto& scene) { scene->template update_state<&Scene::State::event_statistics>( diff --git a/mcp/core/runtime/Plot.hpp b/mcp/core/runtime/Plot.hpp index 48e3fa0..66498ef 100644 --- a/mcp/core/runtime/Plot.hpp +++ b/mcp/core/runtime/Plot.hpp @@ -98,15 +98,10 @@ public: /* 清空旧捕获并请求接下来实际完成的 frame_count 帧 Task DAG。 */ void request_taskflow_trace(std::size_t frame_count); [[nodiscard]] nlohmann::json taskflow_trace() const; - /* 仅捕获 publish 之后真实执行的外接 Task_Graph;不追踪另一张 Plot 的 Render DAG。 */ - void request_post_publish_taskflow_trace(std::size_t frame_count); - [[nodiscard]] nlohmann::json post_publish_taskflow_trace() const; void reset_diagnostics(); private: - friend struct Gallery_Video_Stream; struct Private; void ensure_started(); - void attach_scene_completion(std::unique_ptr completion); std::unique_ptr d; }; } diff --git a/mcp/tests/Control_Path_Benchmarks.cpp b/mcp/tests/Control_Path_Benchmarks.cpp index d5cedb8..d3d6cf7 100644 --- a/mcp/tests/Control_Path_Benchmarks.cpp +++ b/mcp/tests/Control_Path_Benchmarks.cpp @@ -424,6 +424,10 @@ protected: const double wall_ns = runtime_delta("observed_wall_time_ns"); const double busy_ns = runtime_delta("worker_busy_time_ns"); const double cpu_ns = runtime_delta("worker_cpu_time_ns"); + const double cooperative_waits = runtime_delta( + "cooperative_wait_count"); + const double cooperative_wait_ns = runtime_delta( + "cooperative_wait_time_ns"); const double worker_count = runtime_end.result == Tool_Call_Result::ok ? runtime_end.content.at("worker_count").get() : 0.0; state.counters["aggregate_fps"] = elapsed_seconds > 0.0 @@ -436,6 +440,9 @@ protected: ? busy_ns / (wall_ns * worker_count) * 100.0 : 0.0; state.counters["worker_cpu_pct"] = wall_ns > 0.0 && worker_count > 0.0 ? cpu_ns / (wall_ns * worker_count) * 100.0 : 0.0; + state.counters["cooperative_yields"] = cooperative_waits; + state.counters["cooperative_wait_ms"] = + cooperative_wait_ns / 1'000'000.0; state.counters["input_batches"] = static_cast(input_batches); state.counters["input_requests"] = static_cast(input_requests); state.counters["input_request_rate"] = elapsed_seconds > 0.0 @@ -547,7 +554,10 @@ public: << " input_batches=" << counter("input_batches") << " input_requests/s=" << counter("input_request_rate") << " worker_busy=" << counter("worker_busy_pct") << "%" - << " worker_cpu=" << counter("worker_cpu_pct") << "%\n"; + << " worker_cpu=" << counter("worker_cpu_pct") << "%" + << " cooperative_yields=" << counter("cooperative_yields") + << " cooperative_wait=" << counter("cooperative_wait_ms") + << " ms\n"; } output.flags(flags); output.precision(precision); diff --git a/web_server/src/Gallery_Video_Stream.cpp b/web_server/src/Gallery_Video_Stream.cpp index 7d1310a..0b4dda2 100644 --- a/web_server/src/Gallery_Video_Stream.cpp +++ b/web_server/src/Gallery_Video_Stream.cpp @@ -1,4 +1,5 @@ #include "Gallery_Video_Stream.hpp" +#include #include #include #include "frame_sampling/Frame_Sampler.hpp" @@ -57,13 +58,61 @@ void Gallery_Video_Stream::Private::initialize( source_ids.push_back(plot.id); sources.push_back(Source{std::move(plot)}); } - encode_source_slot = sources.size() - 1U; sampler = std::make_unique( gallery_frame_rate, tile_width, tile_height, atlas_columns, - std::move(source_ids), encode_source_slot); + std::move(source_ids), sources.size() - 1U); metric_completion_starts.resize(sources.size()); metric_rendered_starts.resize(sources.size()); } +bool Gallery_Video_Stream::Private::mark_taskflow_trace() { + auto remaining = taskflow_trace_remaining.load(std::memory_order_acquire); + while (remaining != 0) { + if (taskflow_trace_remaining.compare_exchange_weak( + remaining, remaining - 1, std::memory_order_acq_rel, + std::memory_order_acquire)) + return true; + } + return false; +} +void Gallery_Video_Stream::Private::restore_taskflow_trace() noexcept { + taskflow_trace_remaining.fetch_add(1, std::memory_order_release); +} +void Gallery_Video_Stream::Private::store_taskflow_trace( + const Taskflow_Frame_Trace& value) { + const auto trace = std::make_shared( + taskflow_trace_json(value)); + auto state = taskflow_trace_control.load(std::memory_order_acquire); + for (;;) { + const auto requested = static_cast(state >> 32U); + const auto captured = static_cast(state); + if (captured >= requested) return; + taskflow_trace_slots[captured].store(trace, std::memory_order_release); + const auto next = (static_cast(requested) << 32U) | + static_cast(captured + 1U); + if (taskflow_trace_control.compare_exchange_weak( + state, next, std::memory_order_release, + std::memory_order_acquire)) + return; + } +} +nlohmann::json Gallery_Video_Stream::Private::taskflow_trace_response() const { + nlohmann::json frames = nlohmann::json::array(); + const auto state = taskflow_trace_control.load(std::memory_order_acquire); + const auto requested = static_cast(state >> 32U); + const auto captured = static_cast(state); + for (std::uint32_t index = 0; index < captured; ++index) + if (const auto trace = taskflow_trace_slots[index].load( + std::memory_order_acquire)) + frames.push_back(*trace); + const auto remaining = taskflow_trace_remaining.load( + std::memory_order_acquire); + return { + {"protocol", "aethera.taskflow.frames"}, {"version", 1}, + {"requested", requested}, {"remaining", remaining}, + {"captured", frames.size()}, + {"complete", requested != 0 && frames.size() == requested}, + {"frames", std::move(frames)}}; +} bool Gallery_Video_Stream::Private::consumer_accepts( Gallery_Transport_Mode transport) const noexcept { const auto current = consumers.load(std::memory_order_acquire); @@ -86,11 +135,14 @@ bool Gallery_Video_Stream::Private::remove_consumer( replacement->reserve(current->size() - 1); std::ranges::copy_if(*current, std::back_inserter(*replacement), [stream](const Consumer& value) { return value.id != stream; }); + const bool became_empty = replacement->empty(); std::shared_ptr desired{std::move(replacement)}; if (consumers.compare_exchange_weak( current, std::move(desired), std::memory_order_acq_rel, - std::memory_order_acquire)) + std::memory_order_acquire)) { + if (became_empty) sample_timer.cancel(); return true; + } } } void Gallery_Video_Stream::Private::publish( @@ -138,11 +190,7 @@ void Gallery_Video_Stream::Private::accept_frame( if (stopping.load(std::memory_order_acquire) || failed.load(std::memory_order_acquire) || !frame) return; - const auto result = sampler->accept_frame(slot, frame->pixels); - if (slot == encode_source_slot && - result == frame_sampling::Frame_Sampler::Accept_Frame_Result::accepted) { - delivered_clock_ticks.fetch_add(1, std::memory_order_relaxed); - } + static_cast(sampler->accept_frame(slot, frame->pixels)); } void Gallery_Video_Stream::Private::update_metrics( const Plot_Render_Tick& tick, @@ -158,7 +206,7 @@ void Gallery_Video_Stream::Private::update_metrics( metric_encoded_start = encoded_frame_count; metric_pixel_start = pixel_frame_count; metric_sample_start = sampler_state.sampled_frames; - metric_clock_start = delivered_clock_ticks.load(std::memory_order_relaxed); + metric_clock_start = sampling_clock_ticks.load(std::memory_order_relaxed); for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) { metric_completion_starts[slot] = composition.sources[slot].completion_count; @@ -171,7 +219,7 @@ void Gallery_Video_Stream::Private::update_metrics( if (elapsed < metric_interval) return; const double seconds = std::chrono::duration(elapsed).count(); const auto clock_total = - delivered_clock_ticks.load(std::memory_order_relaxed); + sampling_clock_ticks.load(std::memory_order_relaxed); nlohmann::json source_metrics = nlohmann::json::object(); const auto description = sampler->describe(); for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) { @@ -199,10 +247,10 @@ void Gallery_Video_Stream::Private::update_metrics( progress.rendered_correlation_id }, { - "clock_lag_ticks", - progress.rendered_correlation_id != 0 && - tick.sequence >= progress.rendered_correlation_id - ? tick.sequence - progress.rendered_correlation_id + "frame_lag", + progress.rendered_sequence != 0 && + progress.completion_sequence >= progress.rendered_sequence + ? progress.completion_sequence - progress.rendered_sequence : 0 } }; @@ -395,12 +443,11 @@ void Gallery_Video_Stream::bind_plots() { throw std::invalid_argument("gallery video stream requires a Plot"); /* - * 先把媒体尾部作为 encode source Plot 的 post-publish 外接 DAG 安装, - * 再 subscribe/start Plot。这样从第一帧开始 Taskflow trace 都能看到 - * gallery.sample.capture 后分叉到 FFmpeg 与原始像素两个子模型, - * 不存在首帧已经构图后才追加 extension 的初始化窗口。 + * 图集媒体有自己的固定采样时钟,不再依赖某一张 Plot 完成后才能运行。 + * 每个来源只发布 latest;30 FPS tick 到达时若没有新完成帧或上一轮媒体 + * DAG 尚未完成,直接合并,不排队也不占用 Worker 等待。 */ - auto media = std::make_unique("gallery.video.frame"); + auto media = std::make_shared("gallery.video.frame"); auto sample = media->add("gallery.sample.capture", [weak] { if (const auto owner = weak.lock()) { auto& owner_data = static_cast(*owner->d); @@ -469,8 +516,76 @@ void Gallery_Video_Stream::bind_plots() { publish_ffmpeg.precede(pack_pixels); pack_pixels.precede(publish_pixels); publish_pixels.precede(complete); - data.sources[data.encode_source_slot].entry.plot->attach_scene_completion( - std::move(media)); + data.media_graph = media; + data.sample_timer = Frame_Scheduler::instance().make_timer( + [weak](Frame_Scheduler::Tick) { + const auto owner = weak.lock(); + if (!owner) return; + auto& owner_data = static_cast(*owner->d); + if (owner_data.stopping.load(std::memory_order_acquire) || + owner_data.failed.load(std::memory_order_acquire)) + return; + owner_data.sampling_clock_ticks.fetch_add( + 1, std::memory_order_relaxed); + if (!owner_data.sampler->has_pending_frame()) return; + bool idle{}; + if (!owner_data.media_busy.compare_exchange_strong( + idle, true, std::memory_order_acq_rel, + std::memory_order_acquire)) { + owner_data.skipped_sample_ticks.fetch_add( + 1, std::memory_order_relaxed); + return; + } + const auto graph = owner_data.media_graph; + bool trace_reserved = owner_data.mark_taskflow_trace(); + auto trace_frame = trace_reserved + ? std::make_shared(Frame_Identity{ + owner_data.sampling_clock_ticks.load( + std::memory_order_acquire), 0}) + : std::shared_ptr{}; + bool trace_started{}; + if (trace_frame) { + trace_frame->request_taskflow_trace(); + trace_started = aethera::detail::begin_taskflow_trace( + *trace_frame); + if (!trace_started) { + owner_data.restore_taskflow_trace(); + trace_reserved = false; + trace_frame.reset(); + } + } + try { + auto completion = [weak, graph, trace_frame, trace_started] { + static_cast(graph); + if (trace_started) + aethera::detail::finish_taskflow_trace(*trace_frame); + if (const auto completed = weak.lock()) { + auto& completed_data = + static_cast(*completed->d); + if (trace_started) + completed_data.store_taskflow_trace( + trace_frame->take_taskflow_trace()); + completed_data.media_busy.store( + false, std::memory_order_release); + } + }; + if (trace_frame) + aethera::detail::run_taskflow( + *graph, *trace_frame, "gallery.media", + std::move(completion)); + else + aethera::detail::run_taskflow( + *graph, std::move(completion)); + } + catch (...) { + if (trace_started) + aethera::detail::finish_taskflow_trace(*trace_frame); + if (trace_reserved) owner_data.restore_taskflow_trace(); + owner_data.media_busy.store(false, + std::memory_order_release); + owner_data.fail(std::current_exception()); + } + }); for (std::size_t slot = 0; slot < data.sources.size(); ++slot) { auto& source = data.sources[slot]; @@ -512,7 +627,9 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( std::make_shared(std::move(readiness)), std::make_shared(std::move(diagnostics))}; auto current = data.consumers.load(std::memory_order_acquire); + bool activate_sampling{}; for (;;) { + activate_sampling = current->empty(); auto replacement = std::make_shared(); replacement->reserve(current->size() + 1); std::ranges::copy_if(*current, std::back_inserter(*replacement), @@ -527,6 +644,8 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( std::memory_order_acquire)) break; } + if (activate_sampling) + data.sample_timer.start_periodic(gallery_frame_rate); return id; } void Gallery_Video_Stream::unsubscribe(Stream_Id stream) { @@ -542,6 +661,7 @@ void Gallery_Video_Stream::request_video_key_frame() { void Gallery_Video_Stream::shutdown() noexcept { auto& data = static_cast(*d); if (data.stopping.exchange(true, std::memory_order_acq_rel)) return; + data.sample_timer.cancel(); for (auto& source : data.sources) { if (source.stream == 0) continue; try { @@ -609,4 +729,31 @@ nlohmann::json Gallery_Video_Stream::diagnostics() const { } return output; } +void Gallery_Video_Stream::request_taskflow_trace(std::size_t frame_count) { + if (frame_count == 0 || + frame_count > Private::maximum_taskflow_trace_frames) + throw std::invalid_argument( + "Gallery Taskflow trace frame_count must be between 1 and 120"); + auto& data = static_cast(*d); + auto control = data.taskflow_trace_control.load(std::memory_order_acquire); + for (;;) { + const auto requested = static_cast(control >> 32U); + const auto captured = static_cast(control); + if (requested != captured) + throw std::logic_error( + "A Gallery Taskflow trace request is already active"); + const auto next = static_cast(frame_count) << 32U; + if (data.taskflow_trace_control.compare_exchange_weak( + control, next, std::memory_order_release, + std::memory_order_acquire)) + break; + } + for (auto& slot : data.taskflow_trace_slots) + slot.store({}, std::memory_order_release); + data.taskflow_trace_remaining.store(frame_count, + std::memory_order_release); +} +nlohmann::json Gallery_Video_Stream::taskflow_trace() const { + return static_cast(*d).taskflow_trace_response(); +} } diff --git a/web_server/src/Gallery_Video_Stream.hpp b/web_server/src/Gallery_Video_Stream.hpp index 15c2e6d..2b68c04 100644 --- a/web_server/src/Gallery_Video_Stream.hpp +++ b/web_server/src/Gallery_Video_Stream.hpp @@ -57,6 +57,8 @@ struct Gallery_Video_Stream : Def, [[nodiscard]] std::string layout_description( Gallery_Transport_Mode transport) const; [[nodiscard]] nlohmann::json diagnostics() const; + void request_taskflow_trace(std::size_t frame_count); + [[nodiscard]] nlohmann::json taskflow_trace() const; private: void bind_plots(); }; diff --git a/web_server/src/Gallery_Video_Stream.ipp b/web_server/src/Gallery_Video_Stream.ipp index 7b58019..d2fc236 100644 --- a/web_server/src/Gallery_Video_Stream.ipp +++ b/web_server/src/Gallery_Video_Stream.ipp @@ -1,7 +1,9 @@ #pragma once #include "frame_sampling/Frame_Sampler.hpp" +#include #include #include +#include #include #include namespace aethera::web { @@ -28,6 +30,9 @@ struct Gallery_Video_Stream::Private : Prev_Private { Sliding_Statistics publish_ms{600}; /* 两种传输发布统计计算器。 */ std::vector sources{}; /* 已按业务标识排序的稳定图集来源。 */ std::unique_ptr sampler{}; /* latest 图像与 30 FPS deadline 的唯一状态源。 */ + std::shared_ptr media_graph{}; /* 独立媒体时钟运行的持久 DAG;异步完成回调延长其生命周期。 */ + Frame_Scheduler::Timer sample_timer{}; /* 与任意单张 Plot 完成速率解耦的 30 FPS 控制时钟。 */ + std::atomic_bool media_busy{}; /* 单个媒体 DAG 的 latest-only 准入。 */ frame_sampling::ffmpeg::FFmpeg_Frame_Transport ffmpeg_transport; /* FFmpeg 子模型的唯一编码上下文。 */ std::optional active_sample{}; /* 当前 deadline 取得的原生图集视图。 */ std::shared_ptr active_video{}; /* 当前编码结果不可变共享所有权。 */ @@ -37,7 +42,7 @@ struct Gallery_Video_Stream::Private : Prev_Private { std::make_shared()}; /* 不可变订阅集合;写入使用 COW,发布只做一次原子读取。 */ std::atomic_bool stopping{}; /* 关闭开始后拒绝新工作。 */ std::atomic_uint64_t skipped_sample_ticks{}; /* 未到 deadline 或无可写消费者的采样触发数。 */ - std::atomic_uint64_t delivered_clock_ticks{}; /* 已分发给 Plot 的时钟 tick。 */ + std::atomic_uint64_t sampling_clock_ticks{}; /* 有消费者期间触发的独立媒体时钟 tick。 */ std::atomic_bool failed{}; /* 首次 Unknown Failure 后停止热路径。 */ std::optional terminal_failure{}; /* 首次终止失败文本。 */ std::uint64_t encoded_frame_count{}; /* 累计编码成功帧数。 */ @@ -48,8 +53,12 @@ struct Gallery_Video_Stream::Private : Prev_Private { std::uint64_t metric_clock_start{}; /* 窗口起点累计时钟数。 */ std::vector metric_completion_starts{}; /* 各 Plot 窗口起点逻辑完成数。 */ std::vector metric_rendered_starts{}; /* 各 Plot 窗口起点真实画面数。 */ - std::size_t encode_source_slot{}; /* 承载媒体尾部 DAG 的 Plot 槽位。 */ std::uint64_t metric_pixel_start{}; /* 窗口起点累计像素帧数。 */ + static constexpr std::size_t maximum_taskflow_trace_frames{120}; + std::atomic_size_t taskflow_trace_remaining{}; /* 尚待捕获的真实媒体 DAG 次数。 */ + std::atomic_uint64_t taskflow_trace_control{}; /* 高 32 位 requested,低 32 位 captured。 */ + std::array>, + maximum_taskflow_trace_frames> taskflow_trace_slots{}; Private(); ~Private(); template @@ -79,5 +88,9 @@ struct Gallery_Video_Stream::Private : Prev_Private { void publish_ffmpeg_frame(); void publish_pixel_frame(); void complete_media_frame(); + [[nodiscard]] bool mark_taskflow_trace(); + void restore_taskflow_trace() noexcept; + void store_taskflow_trace(const Taskflow_Frame_Trace& trace); + [[nodiscard]] nlohmann::json taskflow_trace_response() const; }; } diff --git a/web_server/src/Gallery_WebSocket.cpp b/web_server/src/Gallery_WebSocket.cpp index ed0f713..61a8715 100644 --- a/web_server/src/Gallery_WebSocket.cpp +++ b/web_server/src/Gallery_WebSocket.cpp @@ -1,6 +1,7 @@ #include "Gallery_WebSocket.hpp" #include "WebRtc_Video_Session.hpp" #include +#include #include #include #include @@ -61,7 +62,17 @@ void Gallery_WebSocket::start() { if (d->attached.exchange(true, std::memory_order_acq_rel)) return; const auto weak = weak_from_this(); if (d->transport == Gallery_Transport_Mode::ffmpeg) { + const auto connection = d->connection.lock(); + if (!connection) + throw std::logic_error( + "WebRTC session requires a live WebSocket connection"); + auto* const owner_loop = + trantor::EventLoop::getEventLoopOfCurrentThread(); + if (!owner_loop) + throw std::logic_error( + "WebRTC session must start on a Drogon I/O thread"); d->video = std::make_unique( + *owner_loop, [weak](std::string signal) { const auto socket = weak.lock(); if (!socket) return; diff --git a/web_server/src/WebRtc_Video_Session.cpp b/web_server/src/WebRtc_Video_Session.cpp index 841738e..a60aebd 100644 --- a/web_server/src/WebRtc_Video_Session.cpp +++ b/web_server/src/WebRtc_Video_Session.cpp @@ -1,7 +1,9 @@ #include "WebRtc_Video_Session.hpp" +#include #include #include #include +#include #include #include #include @@ -35,7 +37,8 @@ void update_maximum(std::atomic_uint64_t& maximum, } } -struct WebRtc_Video_Session::Private { +struct WebRtc_Video_Session::Private final + : std::enable_shared_from_this { struct Callback_State { Signal_Handler signal_handler; /* SDP/ICE 信令的唯一交付出口。 */ Ready_Handler ready_handler; /* Track 打开后请求首个 IDR。 */ @@ -78,13 +81,14 @@ struct WebRtc_Video_Session::Private { std::atomic_size_t transport_buffered_bytes{}; std::atomic_uint64_t readiness_check_count{}; std::atomic_uint64_t readiness_reject_count{}; - std::atomic_uint64_t close_join_started_ns{}; - std::atomic_uint64_t close_join_total_ns{}; - std::atomic_uint64_t close_join_max_ns{}; - std::atomic_bool close_join_active{}; - Private(Signal_Handler signal_handler, Ready_Handler ready_handler, + not_null owner_loop; /* Peer/Track 创建所在 Drogon I/O 线程;服务停止前始终有效。 */ + std::atomic_bool retirement_requested{}; /* close 已撤销入口并要求发送线程退出。 */ + std::atomic_bool retirement_dispatched{}; /* 所属 I/O 线程只接收一次最终资源退役。 */ + Private(trantor::EventLoop& value_owner_loop, + Signal_Handler signal_handler, Ready_Handler ready_handler, Failure_Handler failure_handler) - : callbacks(std::make_shared()) { + : callbacks(std::make_shared()), + owner_loop(&value_owner_loop) { if (!signal_handler || !failure_handler) throw std::invalid_argument( "WebRTC session requires signal and failure handlers"); @@ -93,12 +97,54 @@ struct WebRtc_Video_Session::Private { callbacks->failure_handler = std::move(failure_handler); } + void release_transport() noexcept { + try { + std::lock_guard lock(media_mutex); + const auto track = video_track.exchange({}, std::memory_order_acq_rel); + try { if (track) track->close(); } + catch (...) {} + try { if (peer) peer->close(); } + catch (...) {} + rtp_config.reset(); + peer.reset(); + } + catch (...) {} + } + + void dispatch_retirement() noexcept { + if (!retirement_requested.load(std::memory_order_acquire) || + sender_loop_active.load(std::memory_order_acquire)) + return; + bool expected{}; + if (!retirement_dispatched.compare_exchange_strong( + expected, true, std::memory_order_acq_rel, + std::memory_order_acquire)) + return; + + /* sender_thread 已退出或正处于返回尾声;detach 只解除 std::thread 句柄, + * 不等待。shared ownership 保证 Private 活到所属 I/O 线程完成销毁。 */ + if (sender_thread.joinable()) sender_thread.detach(); + try { + owner_loop->queueInLoop([self = shared_from_this()] { + self->release_transport(); + }); + } + catch (...) { + /* 无法把线程亲和资源交回所属循环属于不可恢复的生命周期故障; + * 禁止在错误线程直接销毁,也不静默回退为同步 join。 */ + std::terminate(); + } + } + void run_sender() noexcept { sender_loop_active.store(true, std::memory_order_release); struct Sender_Exit final { - std::atomic_bool& active; - ~Sender_Exit() { active.store(false, std::memory_order_release); } - } sender_exit{sender_loop_active}; + not_null owner; + ~Sender_Exit() { + owner->sender_loop_active.store(false, std::memory_order_release); + owner->dispatch_retirement(); + } + } sender_exit{this}; for (;;) { Encoded_Video_Frame frame; pending_frames.wait_dequeue(frame); @@ -151,10 +197,11 @@ struct WebRtc_Video_Session::Private { }; WebRtc_Video_Session::WebRtc_Video_Session( - Signal_Handler signal_handler, Ready_Handler ready_handler, + trantor::EventLoop& owner_loop, Signal_Handler signal_handler, + Ready_Handler ready_handler, Failure_Handler failure_handler) - : d(std::make_unique( - std::move(signal_handler), std::move(ready_handler), + : d(std::make_shared( + owner_loop, std::move(signal_handler), std::move(ready_handler), std::move(failure_handler))) {} WebRtc_Video_Session::~WebRtc_Video_Session() { close(); } @@ -254,7 +301,14 @@ void WebRtc_Video_Session::start() { d->peer = std::move(peer); d->video_track.store(std::move(video_track), std::memory_order_release); d->rtp_config = std::move(rtp_config); - d->sender_thread = std::thread([data = d.get()] { data->run_sender(); }); + d->sender_loop_active.store(true, std::memory_order_release); + try { + d->sender_thread = std::thread([data = d] { data->run_sender(); }); + } + catch (...) { + d->sender_loop_active.store(false, std::memory_order_release); + throw; + } } void WebRtc_Video_Session::accept_answer(std::string_view sdp) { @@ -342,8 +396,6 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const { const auto send_total_ns = d->send_total_ns.load(std::memory_order_relaxed); const auto readiness_check_count = d->readiness_check_count.load(std::memory_order_relaxed); - const auto close_join_total_ns = - d->close_join_total_ns.load(std::memory_order_relaxed); const auto milliseconds = [](std::uint64_t nanoseconds) { return static_cast(nanoseconds) / 1'000'000.0; }; @@ -370,11 +422,7 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const { {"current_send_ms", milliseconds(d->send_active.load(std::memory_order_acquire) ? now - d->send_started_ns.load(std::memory_order_relaxed) : 0)}, {"readiness_check_count", readiness_check_count}, - {"readiness_reject_count", d->readiness_reject_count.load(std::memory_order_relaxed)}, - {"close_join_total_ms", milliseconds(close_join_total_ns)}, - {"close_join_max_ms", milliseconds(d->close_join_max_ns.load(std::memory_order_relaxed))}, - {"current_close_join_ms", milliseconds(d->close_join_active.load(std::memory_order_acquire) - ? now - d->close_join_started_ns.load(std::memory_order_relaxed) : 0)}}; + {"readiness_reject_count", d->readiness_reject_count.load(std::memory_order_relaxed)}}; } void WebRtc_Video_Session::close() noexcept { @@ -382,32 +430,12 @@ void WebRtc_Video_Session::close() noexcept { if (d->callbacks->closed.exchange(true, std::memory_order_acq_rel)) return; d->callbacks->track_open.store(false, std::memory_order_release); + d->retirement_requested.store(true, std::memory_order_release); Encoded_Video_Frame discarded; while (d->pending_frames.try_dequeue(discarded)) {} - while (!d->pending_frames.try_enqueue( - Encoded_Video_Frame{})) - std::this_thread::yield(); - if (d->sender_thread.joinable()) { - const auto join_started = steady_nanoseconds(); - d->close_join_started_ns.store(join_started, std::memory_order_relaxed); - d->close_join_active.store(true, std::memory_order_release); - d->sender_thread.join(); - const auto join_elapsed = steady_nanoseconds() - join_started; - d->close_join_active.store(false, std::memory_order_release); - d->close_join_total_ns.fetch_add(join_elapsed, std::memory_order_relaxed); - update_maximum(d->close_join_max_ns, join_elapsed); - } - while (d->pending_frames.try_dequeue(discarded)) {} - d->outstanding_video_frames.store(0, std::memory_order_release); - - std::lock_guard lock(d->media_mutex); - const auto track = d->video_track.exchange({}, std::memory_order_acq_rel); - try { if (track) track->close(); } - catch (...) {} - try { if (d->peer) d->peer->close(); } - catch (...) {} - d->rtp_config.reset(); - d->peer.reset(); + if (!d->pending_frames.try_enqueue(Encoded_Video_Frame{})) + std::terminate(); + d->dispatch_retirement(); } catch (...) {} } diff --git a/web_server/src/WebRtc_Video_Session.hpp b/web_server/src/WebRtc_Video_Session.hpp index e759497..5b9b0cf 100644 --- a/web_server/src/WebRtc_Video_Session.hpp +++ b/web_server/src/WebRtc_Video_Session.hpp @@ -7,6 +7,9 @@ #include #include #include +namespace trantor { +class EventLoop; +} namespace aethera::web { struct WebRtc_Video_Session final { public: @@ -17,7 +20,8 @@ public: using Signal_Handler = std::function; using Ready_Handler = std::function; using Failure_Handler = std::function; - WebRtc_Video_Session(Signal_Handler signal_handler, + WebRtc_Video_Session(trantor::EventLoop& owner_loop, + Signal_Handler signal_handler, Ready_Handler ready_handler, Failure_Handler failure_handler); ~WebRtc_Video_Session(); @@ -33,6 +37,6 @@ public: void close() noexcept; private: struct Private; - std::unique_ptr d; + std::shared_ptr d; }; } diff --git a/web_server/src/Web_Server.cpp b/web_server/src/Web_Server.cpp index a02afd9..426be13 100644 --- a/web_server/src/Web_Server.cpp +++ b/web_server/src/Web_Server.cpp @@ -53,10 +53,8 @@ bool taskflow_graph_contains_gallery_media(const nlohmann::json& graph) { } /* - * Gallery 的媒体 DAG 实际只在组内唯一 encode-source Plot 上执行。诊断接口 - * 在服务端把该真实执行帧中由 Task_Graph 包装器捕获到的媒体 graph 附加到 - * 当前 Plot 的 trace 响应;不复制执行、不伪造 FFmpeg 节点,也不让前端知道 - * encode-source 是哪个 Plot。 + * Gallery 的媒体 DAG 由媒体组自己的固定时钟执行。诊断接口只把该媒体组 + * 实际捕获的 graph 附加到当前 Plot trace;不复制执行,也不伪造节点。 */ nlohmann::json merge_gallery_media_trace(nlohmann::json plot_trace, const nlohmann::json& media_trace) { @@ -134,17 +132,16 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) auto gallery_streams = std::make_shared(); auto plot_media = std::make_shared< std::unordered_map>(); - auto plot_media_source = std::make_shared< + auto plot_media_group = std::make_shared< std::unordered_map>(); - const auto add_media_group = [&gallery_streams, &plot_media, &plot_media_source]( + const auto add_media_group = [&gallery_streams, &plot_media, &plot_media_group]( std::string id, std::vector entries) { if (entries.empty()) return; const auto path = "/ws/gallery?group=" + id; - const auto media_source = entries.back().id; for (const auto& entry : entries) { plot_media->emplace(entry.id, path); - plot_media_source->emplace(entry.id, media_source); + plot_media_group->emplace(entry.id, id); } gallery_streams->emplace( std::move(id), Gallery_Video_Stream::create(std::move(entries))); @@ -237,7 +234,8 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) callback(json_response(output.content)); }, {drogon::Get}); - app.registerHandler("/plot/{1}/taskflow", [control, plot_media_source]( + app.registerHandler("/plot/{1}/taskflow", [control, plot_media_group, + gallery_streams]( const drogon::HttpRequestPtr& request, std::function&& callback, std::string plot_id) { @@ -246,10 +244,16 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) callback(error_response(drogon::k404NotFound, "unknown plot")); return; } - const auto source_found = plot_media_source->find(plot_id); - auto media_plot = source_found == plot_media_source->end() - ? plot : control->find_plot(source_found->second); - if (!media_plot) media_plot = plot; + const auto group_found = plot_media_group->find(plot_id); + const auto stream_found = group_found == plot_media_group->end() + ? gallery_streams->end() + : gallery_streams->find(group_found->second); + if (stream_found == gallery_streams->end()) { + callback(error_response(drogon::k404NotFound, + "plot has no gallery media group")); + return; + } + const auto& media_stream = stream_found->second; try { if (request->method() == drogon::Post) { const auto input = nlohmann::json::parse(request->body()); @@ -265,19 +269,14 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) state.value("captured", std::size_t{}); }; if (trace_active(plot->taskflow_trace()) || - trace_active(media_plot->post_publish_taskflow_trace())) + trace_active(media_stream->taskflow_trace())) throw std::logic_error( "A Taskflow frame trace request is already active"); - /* - * 非 encode-source Plot 只让真实媒体源捕获 post-publish DAG。 - * 不再额外追踪另一张 Plot 的 Render/Prepare/Paint,避免诊断本身 - * 放大 Executor 压力。FFmpeg 仍来自真实 Task_Graph 包装节点。 - */ - media_plot->request_post_publish_taskflow_trace(frame_count); + media_stream->request_taskflow_trace(frame_count); plot->request_taskflow_trace(frame_count); } auto output = plot->taskflow_trace(); - const auto media = media_plot->post_publish_taskflow_trace(); + const auto media = media_stream->taskflow_trace(); output = merge_gallery_media_trace(std::move(output), media); const bool plot_complete = output.value("complete", false); const bool media_complete = media.value("complete", false); diff --git a/web_server/src/frame_sampling/Frame_Sampler.cpp b/web_server/src/frame_sampling/Frame_Sampler.cpp index e81e473..7975887 100644 --- a/web_server/src/frame_sampling/Frame_Sampler.cpp +++ b/web_server/src/frame_sampling/Frame_Sampler.cpp @@ -22,6 +22,7 @@ struct Frame_Sampler::Private { std::size_t sampling_source_slot{}; /* 为样本提供页面媒体时间的固定业务来源。 */ std::atomic_int64_t next_deadline_ns{}; /* 下一次允许生成样本的单调 deadline。 */ std::atomic_uint64_t accepted_source_frames{}; /* 成功发布到 latest 的源帧计数。 */ + std::atomic_uint64_t sampled_source_generation{}; /* 最近样本覆盖到的 accepted_source_frames。 */ std::atomic_uint64_t rejected_source_frames{}; /* 无效或过期源帧计数。 */ std::atomic_uint64_t sampled_frames{}; /* 实际采样计数。 */ std::atomic_uint64_t early_ticks{}; /* deadline 前被合并的触发计数。 */ @@ -63,7 +64,7 @@ Frame_Sampler::Accept_Frame_Result Frame_Sampler::accept_frame( std::size_t slot, std::shared_ptr frame) { const auto result = d->atlas.accept_frame(slot, std::move(frame)); if (result == detail::Gallery_Frame_Atlas::Accept_Frame_Result::accepted) { - d->accepted_source_frames.fetch_add(1, std::memory_order_relaxed); + d->accepted_source_frames.fetch_add(1, std::memory_order_release); return Accept_Frame_Result::accepted; } d->rejected_source_frames.fetch_add(1, std::memory_order_relaxed); @@ -72,10 +73,16 @@ Frame_Sampler::Accept_Frame_Result Frame_Sampler::accept_frame( : Accept_Frame_Result::invalid_frame; } +bool Frame_Sampler::has_pending_frame() const noexcept { + return d->accepted_source_frames.load(std::memory_order_acquire) != + d->sampled_source_generation.load(std::memory_order_acquire); +} + std::optional Frame_Sampler::sample() { const auto now = Clock::now(); const auto now_ns = clock_nanoseconds(now); - const auto compose_sample = [this, now](double deadline_delay_ms) { + const auto compose_sample = [this, now](double deadline_delay_ms, + std::uint64_t generation) { auto composition = d->atlas.compose(); const auto& source = composition.sources[d->sampling_source_slot]; Plot_Render_Tick tick{ @@ -85,9 +92,16 @@ std::optional Frame_Sampler::sample() { 1'000.0, composition.width, composition.height}; + d->sampled_source_generation.store(generation, + std::memory_order_release); return Sampled_Gallery_Frame{ std::move(composition), tick, deadline_delay_ms}; }; + const auto generation = d->accepted_source_frames.load( + std::memory_order_acquire); + if (generation == d->sampled_source_generation.load( + std::memory_order_acquire)) + return std::nullopt; auto deadline_ns = d->next_deadline_ns.load(std::memory_order_acquire); for (;;) { if (deadline_ns == 0) { @@ -97,7 +111,7 @@ std::optional Frame_Sampler::sample() { std::memory_order_acquire)) continue; d->sampled_frames.fetch_add(1, std::memory_order_relaxed); - return compose_sample(0.0); + return compose_sample(0.0, generation); } if (now_ns < deadline_ns) { d->early_ticks.fetch_add(1, std::memory_order_relaxed); @@ -114,7 +128,8 @@ std::optional Frame_Sampler::sample() { continue; d->sampled_frames.fetch_add(1, std::memory_order_relaxed); d->missed_periods.fetch_add(missed, std::memory_order_relaxed); - return compose_sample(static_cast(late_ns) / 1'000'000.0); + return compose_sample(static_cast(late_ns) / 1'000'000.0, + generation); } } diff --git a/web_server/src/frame_sampling/Frame_Sampler.hpp b/web_server/src/frame_sampling/Frame_Sampler.hpp index 2fac74e..5c6669c 100644 --- a/web_server/src/frame_sampling/Frame_Sampler.hpp +++ b/web_server/src/frame_sampling/Frame_Sampler.hpp @@ -50,6 +50,7 @@ public: [[nodiscard]] Accept_Frame_Result accept_frame( std::size_t slot, std::shared_ptr frame); + [[nodiscard]] bool has_pending_frame() const noexcept; [[nodiscard]] std::optional sample(); [[nodiscard]] detail::Gallery_Atlas_Description describe() const; [[nodiscard]] Frame_Sampler_State state() const noexcept; diff --git a/web_server/tests/Frame_Sampler_Tests.cpp b/web_server/tests/Frame_Sampler_Tests.cpp index 85b4443..e0aaccd 100644 --- a/web_server/tests/Frame_Sampler_Tests.cpp +++ b/web_server/tests/Frame_Sampler_Tests.cpp @@ -24,8 +24,10 @@ TEST(Frame_Sampler, First_Deadline_Samples_Latest_Without_Format_Conversion) { Frame_Sampler sampler(30.0, 2, 2, 1, {"plot"}, 0); ASSERT_EQ(sampler.accept_frame(0, native_bgra_frame()), Frame_Sampler::Accept_Frame_Result::accepted); + EXPECT_TRUE(sampler.has_pending_frame()); const auto sample = sampler.sample(); ASSERT_TRUE(sample); + EXPECT_FALSE(sampler.has_pending_frame()); EXPECT_EQ(sample->composition.layout, Plot_Pixel_Layout::bgra8); EXPECT_EQ(sample->composition.pixels[0], std::byte{11}); EXPECT_EQ(sample->composition.pixels[1], std::byte{22}); @@ -37,7 +39,7 @@ TEST(Frame_Sampler, First_Deadline_Samples_Latest_Without_Format_Conversion) { const auto state = sampler.state(); EXPECT_DOUBLE_EQ(state.target_frame_rate_fps, 30.0); EXPECT_EQ(state.sampled_frames, 1U); - EXPECT_EQ(state.early_ticks, 1U); + EXPECT_EQ(state.early_ticks, 0U); } TEST(WebSocket_Pixel_Frame, Header_Describes_Unchanged_Native_Payload) { diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index eb4d0b7..a18dc9b 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -56,7 +56,7 @@ type Gallery_Layout = {kind: "gallery_layout"; protocol: "aethera.gallery.video" tile_width: number; tile_height: number; width: number; height: number; plots: Record}; type Gallery_Source_Metrics = {has_rendered_frame: boolean; logical_completion_rate_fps: number; rendered_frame_rate_fps: number; logical_completion_count: number; rendered_frame_count: number; latest_completion_sequence: number; - latest_rendered_sequence: number; latest_rendered_clock_sequence: number; clock_lag_ticks: number}; + latest_rendered_sequence: number; latest_rendered_clock_sequence: number; frame_lag: number}; type Gallery_Transport_Metrics = {kind: "gallery_metrics"; protocol: "aethera.gallery.video"; version: 2 | 3; clock_sequence: number; encoded_frame_count: number; pixel_frame_count?: number; encoder_backend: "nvenc" | "vulkan_video" | "inactive"; target_frame_rate_fps: number; clock_delivery_rate_fps: number; sampled_frame_rate_fps?: number; @@ -119,17 +119,19 @@ type Taskflow_Frame_Response = {protocol: "aethera.taskflow.frames"; version: 1; captured: number; complete: boolean; frames: Taskflow_Frame_Trace[]; media_requested?: number; media_captured?: number; media_remaining?: number}; type Gallery_Pipeline_State = Record; -type Taskflow_Worker_State = {id: number; task_count: number; current_queue_size: number; current_queue_capacity: number; +type Taskflow_Worker_State = {id: number; task_count: number; entry_queue_size: number; entry_queue_capacity: number; peak_queue_size: number; max_queue_capacity: number; active_task: {native_id: string; type: string; time_ns: number}; - task_time_ns: number; busy_time_ns: number; cpu_time_ns: number; non_cpu_time_ns: number; idle_time_ns: number; + task_time_ns: number; busy_time_ns: number; cpu_time_ns: number; cooperative_wait_count: number; + cooperative_wait_time_ns: number; idle_time_ns: number; min_task_time_ns: number; max_task_time_ns: number; utilization: number; cpu_utilization: number}; type Taskflow_Type_State = {name: string; count: number; total_time_ns: number; min_time_ns: number; max_time_ns: number}; -type Taskflow_Runtime_State = {protocol: "aethera.taskflow.runtime"; version: 1; worker_count: number; active_topologies: number; +type Taskflow_Runtime_State = {protocol: "aethera.taskflow.runtime"; version: 2; worker_count: number; active_topologies: number; active_taskflows: number; peak_active_taskflows: number; completed_taskflows: number; failed_taskflows: number; active_tasks: number; peak_active_tasks: number; active_workers: number; peak_active_workers: number; observed_tasks: number; named_tasks: number; peak_worker_queue_size: number; max_worker_queue_capacity: number; longest_task: {native_id: string; name: string; type: string; time_ns: number}; total_task_time_ns: number; worker_busy_time_ns: number; worker_cpu_time_ns: number; observed_wall_time_ns: number; + cooperative_wait_count: number; cooperative_wait_time_ns: number; worker_utilization: number; worker_cpu_utilization: number; task_types: Taskflow_Type_State[]; workers: Taskflow_Worker_State[]}; @@ -187,9 +189,11 @@ function use_selected_plot_diagnostics(plot: Plot | null) { const sample = async () => { try { const [plot_response, gallery_response] = await Promise.all([ - fetch(plot.diagnostics, {cache: "no-store"}), + fetch(plot.diagnostics, {cache: "no-store", + signal: AbortSignal.timeout(3000)}), gallery_endpoint - ? fetch(gallery_endpoint, {cache: "no-store"}) + ? fetch(gallery_endpoint, {cache: "no-store", + signal: AbortSignal.timeout(3000)}) : Promise.resolve(null) ]); const diagnostics: unknown = await plot_response.json(); @@ -1214,7 +1218,8 @@ function Frame_Policy_Pane({plot, analysis, busy, on_refresh, on_update, on_manu const [capture_busy, set_capture_busy] = useState(false); const [capture_error, set_capture_error] = useState(""); const load_samples = useCallback(async () => { - const request = await fetch(plot.taskflow, {cache: "no-store"}); + const request = await fetch(plot.taskflow, {cache: "no-store", + signal: AbortSignal.timeout(3000)}); const value = await request.json() as Taskflow_Frame_Response & {error?: string}; if (!request.ok) throw new Error(value.error ?? "读取帧策略状态失败"); set_response(value); @@ -1227,9 +1232,15 @@ function Frame_Policy_Pane({plot, analysis, busy, on_refresh, on_update, on_manu }, [load_samples]); useEffect(() => { if (!response || response.requested === 0 || response.complete) return; - const timer = window.setInterval(() => void load_samples().catch(failure => - set_capture_error(failure instanceof Error ? failure.message : "读取帧策略状态失败")), 400); - return () => window.clearInterval(timer); + let stopped = false; + let timer = 0; + const poll = async () => { + try { await load_samples(); } + catch (failure) { set_capture_error(failure instanceof Error ? failure.message : "读取帧策略状态失败"); } + if (!stopped) timer = window.setTimeout(() => void poll(), 400); + }; + void poll(); + return () => { stopped = true; window.clearTimeout(timer); }; }, [load_samples, response?.requested, response?.complete]); const capture = async () => { set_capture_busy(true); set_capture_error(""); set_frame_index(0); @@ -2413,7 +2424,8 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon const [gallery_state, set_gallery_state] = useState(null); const response = scene_response; const read_trace = useCallback(async (endpoint: string) => { - const request = await fetch(endpoint); + const request = await fetch(endpoint, { + signal: AbortSignal.timeout(3000)}); const value = await request.json() as Taskflow_Frame_Response & {error?: string}; if (!request.ok) throw new Error(value.error ?? "读取 Taskflow 帧失败"); return value; @@ -2422,7 +2434,8 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon const group_id = new URL(plot.media, location.href).searchParams.get("group"); const [scene, gallery] = await Promise.all([ read_trace(plot.taskflow), - group_id ? fetch(`/gallery/${encodeURIComponent(group_id)}/diagnostics`, {cache: "no-store"}) + group_id ? fetch(`/gallery/${encodeURIComponent(group_id)}/diagnostics`, { + cache: "no-store", signal: AbortSignal.timeout(3000)}) .then(async request => request.ok ? await request.json() as Gallery_Pipeline_State : null) .catch(() => null) : Promise.resolve(null) ]); @@ -2438,8 +2451,15 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon }, [plot.taskflow]); useEffect(() => { if (!response || response.requested === 0 || response.complete) return; - const timer = window.setInterval(() => void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 帧失败")), 400); - return () => window.clearInterval(timer); + let stopped = false; + let timer = 0; + const poll = async () => { + try { await load(); } + catch (failure) { set_error(failure instanceof Error ? failure.message : "读取 Taskflow 帧失败"); } + if (!stopped) timer = window.setTimeout(() => void poll(), 400); + }; + void poll(); + return () => { stopped = true; window.clearTimeout(timer); }; }, [load, response?.requested, response?.complete]); const capture = async () => { set_busy(true); set_error(""); set_frame_index(0); set_stage_key(""); @@ -2574,33 +2594,42 @@ function Taskflow_Runtime_Pane() { const [state, set_state] = useState(null); const [error, set_error] = useState(""); const load = useCallback(async () => { - const request = await fetch("/taskflow/diagnostics"); + const request = await fetch("/taskflow/diagnostics", { + signal: AbortSignal.timeout(3000)}); const value = await request.json() as Taskflow_Runtime_State & {error?: string}; if (!request.ok) throw new Error(value.error ?? "读取 Taskflow 总体状态失败"); set_state(value); set_error(""); }, []); useEffect(() => { - void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 总体状态失败")); - const timer = window.setInterval(() => void load().catch(failure => set_error(failure instanceof Error ? failure.message : "读取 Taskflow 总体状态失败")), 1000); - return () => window.clearInterval(timer); + let stopped = false; + let timer = 0; + const poll = async () => { + try { await load(); } + catch (failure) { set_error(failure instanceof Error ? failure.message : "读取 Taskflow 总体状态失败"); } + if (!stopped) timer = window.setTimeout(() => void poll(), 1000); + }; + void poll(); + return () => { stopped = true; window.clearTimeout(timer); }; }, [load]); return
全局执行域

Taskflow 总体观测

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

{error ?

{error}

: null}{state ? <>
Worker
{state.worker_count}
任务体墙钟占用
{state.worker_utilization.toFixed(1)}%
-
线程 CPU 占用
{state.worker_cpu_utilization.toFixed(1)}%
+
线程 CPU 占用
{state.worker_cpu_utilization.toFixed(1)}%
活跃 Worker
{state.active_workers}/{state.worker_count}
活跃任务
{state.active_tasks}
活跃 Topology
{state.active_topologies}
活跃 Taskflow
{state.active_taskflows}
累计任务
{state.observed_tasks.toLocaleString("zh-CN")}
失败 Taskflow
{state.failed_taskflows}
队列峰值
{state.peak_worker_queue_size}
最长任务
{nanoseconds(state.longest_task.time_ns)}
+
协作让出
{state.cooperative_wait_count.toLocaleString("zh-CN")} 次 · {nanoseconds(state.cooperative_wait_time_ns)}
-

任务体墙钟占用表示 Worker 正位于任务调用栈中,不等同于 CPU 使用率。线程 CPU 占用按最外层活跃区间累计;两者差值包含 OS 未计费、内核等待和任务体阻塞,但不能仅凭差值断定某一把锁。Executor 拥塞应结合队列峰值和逐帧就绪等待判断。

-
Worker 占用与队列累计值,不在 GET 时重新计算任务样本。
{state.workers.map(worker =>
+

任务体墙钟只累计 Worker 实际执行节点的独占片段,显式协作让出已经扣除。线程 CPU 是 Worker 原生线程从首次任务起的总 CPU,包含 Executor 调度开销,因此不再用“任务墙钟减线程 CPU”伪造非 CPU 等待。队列只报告任务进入时的本地队列样本和累计峰值。

+
Worker 占用、协作让出与本地队列累计值,不在 GET 时重新计算任务样本。
{state.workers.map(worker =>
Worker {worker.id}墙 {worker.utilization.toFixed(1)}% · CPU {worker.cpu_utilization.toFixed(1)}%
-
任务
{worker.task_count}
队列 当前/峰值
{worker.current_queue_size}/{worker.peak_queue_size}
+
任务
{worker.task_count}
进入队列/峰值
{worker.active_task.time_ns ? worker.entry_queue_size : "--"}/{worker.peak_queue_size}
最长
{nanoseconds(worker.max_task_time_ns)}
活跃持续
{worker.active_task.time_ns ? nanoseconds(worker.active_task.time_ns) : "空闲"}
-
累计 CPU
{nanoseconds(worker.cpu_time_ns)}
非 CPU 墙钟
{nanoseconds(worker.non_cpu_time_ns)}
+
线程总 CPU
{nanoseconds(worker.cpu_time_ns)}
独占任务墙钟
{nanoseconds(worker.busy_time_ns)}
+
协作让出
{worker.cooperative_wait_count.toLocaleString("zh-CN")} 次
让出墙钟
{nanoseconds(worker.cooperative_wait_time_ns)}
{worker.active_task.time_ns ? {worker.active_task.type} · {worker.active_task.native_id} : null}
)}
Taskflow 原生任务类型按 Observer TaskType 累计执行次数与耗时。
@@ -2670,7 +2699,7 @@ const Plot_Card = memo(function Plot_Card({plot, selected, policy, gallery, on_p {gallery.pixel_surface ? "Canvas" : "解码"} {gallery.playback.frame_rate_fps.toFixed(1)} FPS 完成 {metrics ? metrics.server_completion_ms.toFixed(1) : "--.-"} ms {encoder_label} {gallery.transport ? (gallery.pixel_surface ? gallery.transport.pixel_pack_average_ms ?? 0 : gallery.transport.encode_average_ms).toFixed(1) : "--.-"} ms - 落后 {source_transport?.has_rendered_frame ? source_transport.clock_lag_ticks : "--"} tick + 画面落后 {source_transport?.has_rendered_frame ? source_transport.frame_lag : "--"} 帧
@@ -2773,11 +2802,11 @@ function Gallery_Grid({plots, selected, policies, layout_scope, galleries, on_po } const workspace_layout_key = "aethera-flexlayout-v7"; -const gallery_transport_key = "aethera-gallery-transport-v1"; +const gallery_transport_key = "aethera-gallery-transport-v2"; function load_gallery_transport(): Gallery_Transport_Mode { - return localStorage.getItem(gallery_transport_key) === "ffmpeg" - ? "ffmpeg" : "websocket_pixels"; + return localStorage.getItem(gallery_transport_key) === "websocket_pixels" + ? "websocket_pixels" : "ffmpeg"; } const layout_labels: Record = { [I18nLabel.Close_Tab]: "关闭标签", @@ -2845,8 +2874,15 @@ function load_workspace_model() { export function App() { const [plots, set_plots] = useState([]); const [category, set_category] = useState("全部"); const [selected, set_selected] = useState(null); const [gallery_transport, set_gallery_transport] = useState(load_gallery_transport); + const visible = useMemo(() => { + if (category === "2D" || category === "3D") + return plots.filter(plot => plot.dimension === category); + return plots; + }, [category, plots]); use_selected_plot_diagnostics(selected); - const gallery_videos = use_gallery_videos(plots, gallery_transport); + /* 只订阅当前维度面板真正消费的 atlas。切到 2D 时不再在后台 + * 同时编码四路 3D atlas,页面筛选本身就是媒体资源的权威需求源。 */ + const gallery_videos = use_gallery_videos(visible, gallery_transport); const [execution_policies, set_execution_policies] = useState({}); const [schema, set_schema] = useState(null); const [schema_busy, set_schema_busy] = useState(false); @@ -2890,16 +2926,12 @@ export function App() { }, []); const reset_frame_diagnostics = () => { if (selected) window.dispatchEvent(new CustomEvent("aethera-reset-frame-diagnostics", {detail: {plot_id: selected.id}})); }; const categories = ["全部", "2D", "3D"]; - const visible = useMemo(() => { - if (category === "2D" || category === "3D") return plots.filter(plot => plot.dimension === category); - return plots; - }, [category, plots]); const gallery =
{selected ? 当前图形 {selected.title} : 点击任意图形后,属性和 DAG 节点状态会自动同步。} + }}>
;