diff --git a/Project_detail_specification.md b/Project_detail_specification.md index 3a2319b..f9bc0e8 100644 --- a/Project_detail_specification.md +++ b/Project_detail_specification.md @@ -18,7 +18,7 @@ * Kernel `Scene` 只负责 2D/3D 共有的 Prepare 数据阶段;Paint、缓存失效、像素合成和异步后端提交由对应渲染模块自己的 Scene、Tag 与 Taskflow 负责。 * `Scene::Private` 是输入事件流的唯一所有者:外部转移事件对象所有权,并发提交到 Def 注册的 MPMC 三缓冲;收集、内部渲染、外部查询分别占用一份队列,Prepare 入口只在交换锁下轮换三者并清空过期查询队列。事件不得独立触发帧。2D 按 Renderable 区域与 Paint 顺序形成接受链,区域默认整个 viewport;3D 无等待提交渲染域,满载时保留当前事件供下一次 Prepare 重试。 * 帧由 Scene 外部创建和持有;每次 `render(frame*)` 只借用该帧并在完成回调返回同一地址。Scene、异步后端和 Web 层只向帧写入固定语义的单调时间点与原始耗时,不保存平均值、分位数或波动等衍生统计。 -* Plot 使用服务端帧时钟调用 `Scene::render(frame*)`;帧策略只决定调用频率,完成回调不得安排下一帧,也不存在浏览器逐帧请求或 request/ack。Scene 回调返回完成帧后,Web 层在 Render Domain 之外把 2D 原生 BGRA 或 3D 原生 RGBA 编码为 H.264,经 LibDataChannel WebRTC 视频轨直接发布;WebSocket 只承载 SDP/ICE、输入事件和诊断 JSON。多个订阅者共享同一编码结果,关闭视频时 3D 必须使用 diagnostics 输出并跳过 GPU 像素回读。 +* Plot 使用服务端帧时钟调用 `Scene::render(frame*)`;`fixed_rate` 只由 Frame_Scheduler 周期时钟驱动,完成回调不得改变其 deadline。`maximum_rate` 是唯一允许完成回调在释放 Plot 准入后异步投递下一帧请求的特殊模式;回调不得直接重入 `render()`、不得同步等待,并且物理帧槽耗尽时只能由帧退役事件解除背压。系统不存在浏览器逐帧 request/ack。Scene 回调返回完成帧后,Web 层在 Render Domain 之外把 2D 原生 BGRA 或 3D 原生 RGBA 编码为 H.264,经 LibDataChannel WebRTC 视频轨直接发布;WebSocket 只承载 SDP/ICE、输入事件和诊断 JSON。多个订阅者共享同一编码结果,关闭视频时 3D 必须使用 diagnostics 输出并跳过 GPU 像素回读。 * 相机是 3D Scene 组件,只允许定义在 `render_3D/camera`;Kernel 和 2D 不得依赖相机类型。Web 层只为已有 `Camera_3D` 增加协议描述,不复制相机配置。Gallery 服务同一时刻只允许一个页面实例持有;该页面的所有 Plot 连接共享页面令牌,其他标签页或浏览器实例必须被拒绝。 * `Time_Axis` 是时间与 tick 的唯一权威来源;使用层先推进时间轴,再把同一 tick 分发给所有相关数据图元。时间窗口从第一条数据起始终锚定最新 tick,未产生数据的槽位保持背景。 * 图表选区保存两根轴上的数据范围,绘制时才映射为像素;选区作为独立 Renderable 在使用层与图元组合,禁止在各图元内复制选区状态。 diff --git a/kernel/src/kernel/Frame_Pacing_Policy.cpp b/kernel/src/kernel/Frame_Pacing_Policy.cpp index 56515a9..6735f01 100644 --- a/kernel/src/kernel/Frame_Pacing_Policy.cpp +++ b/kernel/src/kernel/Frame_Pacing_Policy.cpp @@ -1,81 +1,39 @@ #include "Frame_Pacing_Policy.hpp" #include #include -#include namespace aethera { -Frame_Pacing_Policy::Frame_Pacing_Policy() - : properties_(std::make_shared()) {} +Frame_Pacing_Policy::Frame_Pacing_Policy() = default; Frame_Pacing_Properties Frame_Pacing_Policy::read() const { - return *properties_.load(std::memory_order_acquire); -} - -template -void Frame_Pacing_Policy::update(Edit&& edit) { - auto current = properties_.load(std::memory_order_acquire); - for (;;) { - auto next = std::make_shared(*current); - std::forward(edit)(*next); - std::shared_ptr desired = next; - if (properties_.compare_exchange_weak( - current, desired, std::memory_order_release, - std::memory_order_acquire)) - return; - } + return { + render_enabled_.load(std::memory_order_acquire), + video_enabled_.load(std::memory_order_acquire), + mode_.load(std::memory_order_acquire), + fixed_rate_fps_.load(std::memory_order_acquire)}; } void Frame_Pacing_Policy::set_render_enabled(bool enabled) { - update([enabled](auto& value) { value.render_enabled = enabled; }); + render_enabled_.store(enabled, std::memory_order_release); } void Frame_Pacing_Policy::set_video_enabled(bool enabled) { - update([enabled](auto& value) { value.video_enabled = enabled; }); + video_enabled_.store(enabled, std::memory_order_release); } void Frame_Pacing_Policy::set_mode(Frame_Pacing_Mode mode) { - update([mode](auto& value) { value.mode = mode; }); + mode_.store(mode, std::memory_order_release); } void Frame_Pacing_Policy::set_fixed_rate(double fps) { if (!std::isfinite(fps) || fps <= 0.0) throw std::invalid_argument("frame pacing fps must be positive"); - update([fps](auto& value) { value.fixed_rate_fps = fps; }); + fixed_rate_fps_.store(fps, std::memory_order_release); } -double Frame_Pacing_Policy::scheduled_rate_fps() const { - const auto value = read(); - if (!value.render_enabled || value.mode == Frame_Pacing_Mode::manual) - return 0.0; - return value.mode == Frame_Pacing_Mode::maximum_rate - ? value.maximum_rate_fps : value.fixed_rate_fps; +bool Frame_Pacing_Policy::request_immediate() const noexcept { + return render_enabled_.load(std::memory_order_acquire); } -bool Frame_Pacing_Policy::accept_periodic_tick(double time_milliseconds) { - static_cast(time_milliseconds); - const auto value = read(); - if (!value.render_enabled || value.mode == Frame_Pacing_Mode::manual) - return false; - if (skip_next_periodic_.exchange(false, std::memory_order_acq_rel)) - return false; - return true; -} - -bool Frame_Pacing_Policy::request_immediate() { - const auto value = read(); - if (!value.render_enabled) return false; - /* immediate 与下一次周期机会合并,真正的 busy/latest 语义由 Plot 管理。 */ - skip_next_periodic_.store(true, std::memory_order_release); - return true; -} - -void Frame_Pacing_Policy::frame_submitted() {} - -bool Frame_Pacing_Policy::frame_completed() { - return false; -} - -void Frame_Pacing_Policy::frame_rejected() {} - } diff --git a/kernel/src/kernel/Frame_Pacing_Policy.hpp b/kernel/src/kernel/Frame_Pacing_Policy.hpp index fc3731a..fd26b20 100644 --- a/kernel/src/kernel/Frame_Pacing_Policy.hpp +++ b/kernel/src/kernel/Frame_Pacing_Policy.hpp @@ -1,7 +1,6 @@ #pragma once #include #include -#include namespace aethera { @@ -12,11 +11,10 @@ enum struct Frame_Pacing_Mode : std::uint8_t { }; struct Frame_Pacing_Properties { - bool render_enabled{true}; - bool video_enabled{true}; - Frame_Pacing_Mode mode{Frame_Pacing_Mode::fixed_rate}; - double fixed_rate_fps{30.0}; - double maximum_rate_fps{100.0}; + bool render_enabled{true}; /* 查询时 Plot 是否接受任何帧请求。 */ + bool video_enabled{true}; /* 查询时完成帧是否需要生成视频像素。 */ + Frame_Pacing_Mode mode{Frame_Pacing_Mode::fixed_rate}; /* 查询时服务端帧驱动模式。 */ + double fixed_rate_fps{30.0}; /* fixed_rate 模式目标帧率。 */ }; /* @@ -32,24 +30,13 @@ struct Frame_Pacing_Policy final { void set_mode(Frame_Pacing_Mode mode); void set_fixed_rate(double fps); - /* 当前策略应向全局 Frame_Scheduler 注册的频率;0 表示事件驱动/关闭周期。 */ - [[nodiscard]] double scheduled_rate_fps() const; - - /* - * 周期 tick 只判断策略/周期语义,不再承担 Scene in-flight 门控。 - * 实际渲染准入由调用方在 Scene::advance 前控制;忙时保留 latest pending tick。 - */ - [[nodiscard]] bool accept_periodic_tick(double time_milliseconds); - [[nodiscard]] bool request_immediate(); - void frame_submitted(); - [[nodiscard]] bool frame_completed(); - void frame_rejected(); + [[nodiscard]] bool request_immediate() const noexcept; private: - template void update(Edit&& edit); - - std::atomic> properties_; - std::atomic_bool skip_next_periodic_{}; + std::atomic_bool render_enabled_{true}; /* 当前 Plot 是否接受任何帧请求。 */ + std::atomic_bool video_enabled_{true}; /* 完成帧是否需要生成视频像素。 */ + std::atomic mode_{Frame_Pacing_Mode::fixed_rate}; /* 当前服务端帧驱动模式。 */ + std::atomic fixed_rate_fps_{30.0}; /* fixed_rate 模式目标帧率。 */ }; } diff --git a/kernel/src/kernel/Frame_Scheduler.cpp b/kernel/src/kernel/Frame_Scheduler.cpp index f0c96ae..d4eab34 100644 --- a/kernel/src/kernel/Frame_Scheduler.cpp +++ b/kernel/src/kernel/Frame_Scheduler.cpp @@ -4,7 +4,6 @@ #include #include #include -#include #include #include #include @@ -63,7 +62,6 @@ struct Frame_Scheduler::Private { Command_Type type{}; std::shared_ptr state{}; Nanoseconds value{}; - std::shared_ptr> completion{}; }; TimerWheel wheel{}; @@ -72,7 +70,6 @@ struct Frame_Scheduler::Private { moodycamel::BlockingConcurrentQueue commands{256}; std::atomic_uint64_t next_id{1}; std::atomic_bool stopping{}; - std::thread::id control_thread_id{}; std::jthread thread{}; Private() : thread([this](std::stop_token stop) { run(stop); }) {} @@ -149,10 +146,6 @@ struct Frame_Scheduler::Private { case Command_Type::stop: break; } - if (command.completion) { - try { command.completion->set_value(); } - catch (...) {} - } } void process_commands(Scheduler_Clock::time_point now) { @@ -210,7 +203,6 @@ struct Frame_Scheduler::Private { } void run(std::stop_token stop) noexcept { - control_thread_id = std::this_thread::get_id(); try { while (!stop.stop_requested() && !stopping.load(std::memory_order_acquire)) { @@ -301,25 +293,6 @@ void Frame_Scheduler::Timer::cancel() noexcept { state_->owner->enqueue({Private::Command_Type::cancel, state_, {}}); } - -void Frame_Scheduler::Timer::cancel_and_wait() noexcept { - if (!state_ || !state_->owner) return; - state_->alive.store(false, std::memory_order_release); - auto* owner = state_->owner; - if (owner->control_thread_id == std::this_thread::get_id()) { - state_->event.cancel(); - state_->mode = State::Mode::stopped; - return; - } - try { - auto completion = std::make_shared>(); - auto future = completion->get_future(); - owner->enqueue({Private::Command_Type::cancel, state_, {}, completion}); - future.wait(); - } - catch (...) {} -} - bool Frame_Scheduler::Timer::valid() const noexcept { return static_cast(state_); } diff --git a/kernel/src/kernel/Frame_Scheduler.hpp b/kernel/src/kernel/Frame_Scheduler.hpp index 27d5054..5d09301 100644 --- a/kernel/src/kernel/Frame_Scheduler.hpp +++ b/kernel/src/kernel/Frame_Scheduler.hpp @@ -40,7 +40,6 @@ struct Frame_Scheduler final { /* 单次延迟;delay <= 0 会在下一个 Scheduler 周期尽快触发。 */ void start_once(std::chrono::nanoseconds delay); void cancel() noexcept; - void cancel_and_wait() noexcept; /* 销毁拥有 callback 的对象前使用,保证控制线程不再执行该 callback。 */ [[nodiscard]] bool valid() const noexcept; private: diff --git a/kernel/src/kernel/frame.cpp b/kernel/src/kernel/frame.cpp index 9b90628..265fada 100644 --- a/kernel/src/kernel/frame.cpp +++ b/kernel/src/kernel/frame.cpp @@ -24,6 +24,7 @@ struct Render_Frame::Private { std::vector tasks{}; /* 仅对应 worker 写入,捕获结束后统一读取。 */ }; Frame_Identity identity{}; /* 外部帧管理器提供且终生不变的帧身份。 */ + Frame_Request_Source request_source{}; /* 当前逻辑帧的策略触发来源。 */ std::chrono::steady_clock::time_point created_at{}; /* 所有 elapsed_ns 使用的单调时钟原点。 */ std::uint64_t created_time_unix_ns{}; /* 用于跨进程展示的创建 Unix 时间,单位为纳秒。 */ std::array markers{}; /* 每种时间点首次出现时的 elapsed_ns 加一编码。 */ @@ -35,8 +36,10 @@ struct Render_Frame::Private { Render_Frame::Render_Frame(Frame_Identity identity) : d(std::make_unique()) { begin(identity); } -void Render_Frame::begin(Frame_Identity identity) noexcept { +void Render_Frame::begin(Frame_Identity identity, + Frame_Request_Source source) noexcept { d->identity = identity; + d->request_source = source; d->created_at = std::chrono::steady_clock::now(); const auto system_now = std::chrono::system_clock::now().time_since_epoch(); d->created_time_unix_ns = static_cast(std::chrono::duration_cast(system_now).count()); @@ -154,6 +157,7 @@ bool Render_Frame::taskflow_trace_requested() const noexcept { Taskflow_Frame_Trace Render_Frame::take_taskflow_trace() { Taskflow_Frame_Trace result{}; result.identity = d->identity; + result.request_source = d->request_source; result.created_time_unix_ns = d->created_time_unix_ns; result.worker_count = d->taskflow_workers.size(); for (std::size_t index = 0; index < marker_count; ++index) { diff --git a/kernel/src/kernel/frame.hpp b/kernel/src/kernel/frame.hpp index 92a0695..e518f8b 100644 --- a/kernel/src/kernel/frame.hpp +++ b/kernel/src/kernel/frame.hpp @@ -9,6 +9,12 @@ namespace detail { struct Taskflow_Frame_Access; } struct Frame_Statistics_Sample; +enum struct Frame_Request_Source : std::uint8_t { + unspecified, + periodic, + immediate, + maximum_rate +}; enum struct Frame_Trace_Marker : std::uint8_t { created, plot_update_started, @@ -121,6 +127,7 @@ struct Taskflow_Task_Trace { }; struct Taskflow_Frame_Trace { Frame_Identity identity{}; /* 该执行图所属逻辑渲染帧。 */ + Frame_Request_Source request_source{}; /* 该逻辑帧由周期、手动或最大吞吐策略触发。 */ std::uint64_t created_time_unix_ns{}; /* 帧创建 Unix 时间,单位纳秒。 */ std::size_t worker_count{}; /* 捕获时全局 Executor 的 worker 数。 */ std::vector markers{}; /* 与该帧 DAG 共用时间原点的原始流水线时间点。 */ @@ -142,7 +149,8 @@ public: [[nodiscard]] Taskflow_Frame_Trace take_taskflow_trace(); protected: /* 物理帧槽再次承载新逻辑帧时,重建其唯一身份和诊断时间原点。 */ - void begin(Frame_Identity identity) noexcept; + void begin(Frame_Identity identity, + Frame_Request_Source source = Frame_Request_Source::unspecified) noexcept; private: struct Private; friend struct detail::Taskflow_Frame_Access; diff --git a/kernel/src/test/render_test.cpp b/kernel/src/test/render_test.cpp index d237fb4..e4150e3 100644 --- a/kernel/src/test/render_test.cpp +++ b/kernel/src/test/render_test.cpp @@ -1,4 +1,5 @@ #include "scene.hpp" +#include "Frame_Pacing_Policy.hpp" #include #include #include @@ -50,6 +51,24 @@ using Scene = aethera::Scene; inline Direct_Renderable::Private& Direct_Renderable::data_for_test() { return static_cast(*d); } + +TEST(frame_pacing_policy, independent_properties_do_not_publish_copied_versions) { + aethera::Frame_Pacing_Policy policy; + policy.set_mode(aethera::Frame_Pacing_Mode::maximum_rate); + policy.set_fixed_rate(47.5); + policy.set_video_enabled(false); + + const auto properties = policy.read(); + EXPECT_TRUE(properties.render_enabled); + EXPECT_FALSE(properties.video_enabled); + EXPECT_EQ(properties.mode, aethera::Frame_Pacing_Mode::maximum_rate); + EXPECT_DOUBLE_EQ(properties.fixed_rate_fps, 47.5); + EXPECT_TRUE(policy.request_immediate()); + + policy.set_render_enabled(false); + EXPECT_FALSE(policy.request_immediate()); + EXPECT_THROW(policy.set_fixed_rate(0.0), std::invalid_argument); +} inline Graph_Renderable::Private& Graph_Renderable::data_for_test() { return static_cast(*d); } diff --git a/render_2D/render_2D/base/Frame_2D.cpp b/render_2D/render_2D/base/Frame_2D.cpp index 9f44c35..570989b 100644 --- a/render_2D/render_2D/base/Frame_2D.cpp +++ b/render_2D/render_2D/base/Frame_2D.cpp @@ -10,11 +10,12 @@ Frame_2D::Frame_2D(Frame_Identity identity, Pixel_Format output_format) : Render_Frame(identity), d(std::make_unique()) { begin(identity, output_format); } -void Frame_2D::begin(Frame_Identity identity, Pixel_Format output_format) { +void Frame_2D::begin(Frame_Identity identity, Pixel_Format output_format, + Frame_Request_Source source) { if (std::ranges::find(supported_pixel_formats, output_format) == supported_pixel_formats.end()) throw std::invalid_argument("Frame_2D output pixel format is unsupported"); - Render_Frame::begin(identity); + Render_Frame::begin(identity, source); d->output_format = output_format; } Frame_2D::~Frame_2D() = default; diff --git a/render_2D/render_2D/base/Frame_2D.hpp b/render_2D/render_2D/base/Frame_2D.hpp index 52aaae5..4c55ad9 100644 --- a/render_2D/render_2D/base/Frame_2D.hpp +++ b/render_2D/render_2D/base/Frame_2D.hpp @@ -23,7 +23,8 @@ public: ~Frame_2D() override; /* 保留二维后端图像分配,只重启物理槽承载的新逻辑帧。 */ void begin(Frame_Identity identity, - Pixel_Format output_format = native_pixel_format); + Pixel_Format output_format = native_pixel_format, + Frame_Request_Source source = Frame_Request_Source::unspecified); [[nodiscard]] Pixel_Format output_format() const noexcept; [[nodiscard]] Image_View image() const; /* 按本帧声明的输出格式生成连续像素;原生格式只做必要的行打包。 */ diff --git a/render_3D/render_3D/base/Frame_3D.cpp b/render_3D/render_3D/base/Frame_3D.cpp index 77cd87a..bc66fff 100644 --- a/render_3D/render_3D/base/Frame_3D.cpp +++ b/render_3D/render_3D/base/Frame_3D.cpp @@ -14,10 +14,10 @@ Frame_3D::Frame_3D(Frame_Identity identity, Frame_3D_Output output, begin(identity, output, output_format); } void Frame_3D::begin(Frame_Identity identity, Frame_3D_Output output, - Pixel_Format output_format) { + Pixel_Format output_format, Frame_Request_Source source) { if (output_format != native_pixel_format) throw std::invalid_argument("Frame_3D output pixel format is unsupported"); - Render_Frame::begin(identity); + Render_Frame::begin(identity, source); d->output = output; d->output_format = output_format; d->extent = {}; diff --git a/render_3D/render_3D/base/Frame_3D.hpp b/render_3D/render_3D/base/Frame_3D.hpp index 814fc1b..8c77c07 100644 --- a/render_3D/render_3D/base/Frame_3D.hpp +++ b/render_3D/render_3D/base/Frame_3D.hpp @@ -30,7 +30,8 @@ public: /* 释放上一逻辑帧的共享读回引用,并重启同一物理帧槽。 */ void begin(Frame_Identity identity, Frame_3D_Output output = Frame_3D_Output::pixels, - Pixel_Format output_format = native_pixel_format); + Pixel_Format output_format = native_pixel_format, + Frame_Request_Source source = Frame_Request_Source::unspecified); [[nodiscard]] Extent extent() const noexcept; [[nodiscard]] std::span pixels() const noexcept; [[nodiscard]] std::shared_ptr> share_pixels() const noexcept; diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp index 647c9fc..ae6cc91 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp @@ -22,21 +22,23 @@ void update_peak(Value& peak, Value value) noexcept { Gpu_Completion_Service::Private::Private() = default; -Gpu_Completion_Service::Private::~Private() { - stopping.store(true, std::memory_order_release); - poll_timer.cancel_and_wait(); -} +Gpu_Completion_Service::Private::~Private() = default; Gpu_Completion_Service::Gpu_Completion_Service() = default; Gpu_Completion_Service::~Gpu_Completion_Service() = default; Gpu_Completion_Service& Gpu_Completion_Service::instance() { - static auto service = [] { + /* + * GPU completion 域与进程同寿命。主动泄放所有权可避免静态析构阶段 + * 为等待 Frame_Scheduler callback 退出而引入同步取消;进程终止负责 + * 回收其内存,Scheduler 仍按自身 RAII 顺序停止控制线程。 + */ + static auto* service = [] { auto built = Gpu_Completion_Service::Builder{}.build(); if (!built) throw std::logic_error( "GPU completion service dependency graph is invalid"); - return std::move(*built); + return built->release(); }(); return *service; } diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index 6777d49..e053796 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -3,6 +3,7 @@ #include "Taskflow_Trace_Json.hpp" #include #include +#include #include #include #include @@ -51,24 +52,6 @@ std::string exception_description(const std::exception_ptr& failure) { return "empty Plot failure"; } -struct Frame_Policy final { -public: - [[nodiscard]] Frame_Pacing_Properties read() const { return pacing.read(); } - [[nodiscard]] nlohmann::json schema() const; - [[nodiscard]] nlohmann::json write_prop(std::string_view key, - const nlohmann::json& value); - [[nodiscard]] bool accept_periodic_tick(double time_milliseconds) { - return pacing.accept_periodic_tick(time_milliseconds); - } - [[nodiscard]] bool request_immediate() { return pacing.request_immediate(); } - [[nodiscard]] double scheduled_rate_fps() const { return pacing.scheduled_rate_fps(); } - void frame_submitted() { pacing.frame_submitted(); } - [[nodiscard]] bool frame_completed() { return pacing.frame_completed(); } - void frame_rejected() { pacing.frame_rejected(); } -private: - Frame_Pacing_Policy pacing{}; -}; - std::string_view pacing_mode_name(Frame_Pacing_Mode mode) { const auto name = magic_enum::enum_name(mode); if (name.empty()) throw std::logic_error("unknown frame pacing mode"); @@ -93,8 +76,8 @@ std::string_view pixel_format_name(render_3d::Pixel_Format format) { throw std::logic_error("unknown 3D pixel format"); } -nlohmann::json Frame_Policy::schema() const { - const auto current = read(); +nlohmann::json frame_policy_schema(const Frame_Pacing_Policy& pacing) { + const auto current = pacing.read(); return { {"id", "frame-analysis"}, {"label", "渲染与媒体流水线"}, {"kind", "analysis"}, {"fields", nlohmann::json::array({ @@ -113,19 +96,20 @@ nlohmann::json Frame_Policy::schema() const { {"options", nlohmann::json::array({ {{"value", "manual"}, {"label", "手动渲染"}}, {{"value", "fixed_rate"}, {"label", "固定频率"}}, - {{"value", "maximum_rate"}, {"label", "最大频率"}} + {{"value", "maximum_rate"}, {"label", "最大吞吐"}} })}}, {{"key", "fixed_rate_fps"}, {"label", "目标帧率"}, {"editor", "number"}, {"editable", true}, {"minimum", 0.1}, {"maximum", 100.0}, {"step", 0.1}, - {"description", "当前 Scene 独立目标帧率;周期策略保存在 Kernel Frame_Pacing_Policy。"}, + {"description", "仅 fixed_rate 使用;maximum_rate 在完成准入释放后异步自驱下一帧。"}, {"technical_description", "Independent per-scene target frame rate."}, {"value", current.fixed_rate_fps}} })} }; } -nlohmann::json Frame_Policy::write_prop(std::string_view key, - const nlohmann::json& value) { +nlohmann::json write_frame_policy_prop(Frame_Pacing_Policy& pacing, + std::string_view key, + const nlohmann::json& value) { if (key == "render_enabled" || key == "video_enabled") { if (!value.is_boolean()) return {{"success", false}, {"error", "frame policy switch requires a boolean"}}; @@ -275,6 +259,7 @@ nlohmann::json taskflow_trace_json( return { {"sequence", trace.identity.sequence}, {"correlation_id", trace.identity.correlation_id}, + {"request_source", magic_enum::enum_name(trace.request_source)}, {"created_time_unix_ns", trace.created_time_unix_ns}, {"worker_count", trace.worker_count}, {"markers", std::move(markers)}, @@ -344,6 +329,12 @@ struct Plot::Private { consuming /* 外接 Taskflow 正在消费已发布帧;允许下一帧渲染。 */ }; + enum struct Render_Admission_State : std::uint8_t { + ready, + rendering, + frame_slots_exhausted + }; + struct Managed_Frame { std::chrono::microseconds presentation_time{}; /* 共享页面时钟产生的媒体时间戳。 */ Frame frame{}; /* 三缓冲物理槽拥有且反复承载逻辑帧。 */ @@ -375,7 +366,7 @@ struct Plot::Private { std::atomic_uint64_t next_stream_id{1}; std::atomic> terminal_failure{}; /* 首次 Plot Unknown Failure 的唯一终止状态。 */ std::uint64_t next_frame_sequence{1}; - Frame_Policy frame_policy{}; + Frame_Pacing_Policy frame_policy{}; Frame_Scheduler::Timer frame_timer{}; /* 每 Plot/Scene 只有轻量时间轮节点,不持有线程。 */ static constexpr std::size_t scene_frame_capacity{3}; std::array frame_slots{}; /* Scene 与外接消费者共享生命周期的稳定三缓冲。 */ @@ -383,15 +374,15 @@ struct Plot::Private { std::atomic retired_frames{}; /* 多回调生产、唯一 Taskflow 任务消费。 */ std::atomic_bool retired_frame_task_scheduled{}; Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */ - std::atomic> pending_tick{}; + moodycamel::ConcurrentQueue tick_requests{}; /* 多生产者提交、唯一短任务消费的帧请求流。 */ + std::optional deferred_tick{}; /* 仅 tick consumer 任务访问的 latest 延后请求。 */ + std::atomic_uint64_t tick_request_generation{}; /* 请求入队后推进,用于关闭 consumer 尾部唤醒竞争窗口。 */ std::atomic_bool tick_task_scheduled{}; /* 唯一短任务准入;不占用 Worker 等待。 */ - std::atomic_bool render_admission_busy{}; /* view->update/Scene::advance 到 pixel publish 的唯一准入门。 */ + std::atomic render_admission{Render_Admission_State::ready}; /* Plot 渲染准入及物理槽背压的唯一状态源。 */ std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()}; std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */ std::atomic_uint64_t coalesced_tick_count{}; /* 尚未消费时被更新 tick 替换的旧 tick 总数。 */ std::atomic_uint64_t policy_skip_count{}; /* 帧策略拒绝的 tick 总数。 */ - std::atomic_uint64_t preparation_busy_count{}; /* Render admission 忙时被合并为 latest pending 的 tick。 */ - std::atomic_uint64_t deferred_resume_count{}; /* pixel publish 后立即唤醒 latest pending 的次数。 */ std::atomic_uint64_t frame_slot_busy_count{}; /* 三个物理帧槽均被占用的提交次数。 */ std::atomic_uint64_t scene_rejection_count{}; /* Scene 单帧准入拒绝的提交次数。 */ std::atomic_uint64_t submitted_frame_count{}; /* 成功提交给 Scene 的帧总数。 */ @@ -440,7 +431,8 @@ struct Plot::Private { [[nodiscard]] nlohmann::json schema() const; [[nodiscard]] Stream_Snapshot stream_snapshot() const; void publish(std::shared_ptr frame) noexcept; - void defer_tick(const Plot_Render_Tick& tick); + void submit_tick_request(Plot_Render_Tick tick); + void keep_latest_tick(const Plot_Render_Tick& tick); void arm_tick_consumer(std::weak_ptr lifetime); void release_render_admission(std::weak_ptr lifetime); void consume_tick(std::weak_ptr lifetime); @@ -491,7 +483,7 @@ void Plot::Private::fail(std::exception_ptr failure) noexcept { nlohmann::json Plot::Private::schema() const { auto result = view->schema(); - auto analysis = frame_policy.schema(); + auto analysis = frame_policy_schema(frame_policy); const auto generator = view->data_generator_schema(); if (!generator.is_null()) analysis["data_generator"] = generator; result["frame_analysis"] = std::move(analysis); @@ -545,39 +537,55 @@ void Plot::Private::publish( void Plot::Private::refresh_schedule() { if (!frame_timer.valid()) return; const auto current_consumers = consumers.load(std::memory_order_acquire); - const double fps = frame_policy.scheduled_rate_fps(); - if (fps <= 0.0 || current_consumers->empty()) frame_timer.cancel(); - else frame_timer.start_periodic(fps); + const auto pacing = frame_policy.read(); + if (!pacing.render_enabled || current_consumers->empty()) { + frame_timer.cancel(); + return; + } + if (pacing.mode == Frame_Pacing_Mode::fixed_rate) { + frame_timer.start_periodic(pacing.fixed_rate_fps); + return; + } + frame_timer.cancel(); + if (pacing.mode != Frame_Pacing_Mode::maximum_rate) return; + const auto now = std::chrono::steady_clock::now(); + submit_tick_request(Plot_Render_Tick{ + .issued_at = now, + .time_milliseconds = std::chrono::duration( + now - clock_origin).count(), + .source = Plot_Render_Tick_Source::maximum_rate}); + arm_tick_consumer(lifetime); } +void Plot::Private::submit_tick_request(Plot_Render_Tick tick) { + if (!tick_requests.enqueue(std::move(tick))) + throw std::bad_alloc{}; + tick_request_generation.fetch_add(1, std::memory_order_release); +} -void Plot::Private::defer_tick(const Plot_Render_Tick& tick) { - const auto next = std::make_shared(tick); - auto current = pending_tick.load(std::memory_order_acquire); - for (;;) { - if (current) { - const bool current_immediate = current->sequence == 0; - const bool next_immediate = tick.sequence == 0; - if ((current_immediate && !next_immediate) || - (current_immediate == next_immediate && - current->issued_at >= tick.issued_at)) - return; +void Plot::Private::keep_latest_tick(const Plot_Render_Tick& tick) { + const auto priority = [](Plot_Render_Tick_Source source) { + switch (source) { + case Plot_Render_Tick_Source::periodic: return 0; + case Plot_Render_Tick_Source::maximum_rate: return 1; + case Plot_Render_Tick_Source::immediate: return 2; } - if (pending_tick.compare_exchange_weak( - current, next, std::memory_order_acq_rel, - std::memory_order_acquire)) { - if (current) - coalesced_tick_count.fetch_add(1, std::memory_order_relaxed); + return 0; + }; + if (deferred_tick) { + const auto current_priority = priority(deferred_tick->source); + const auto next_priority = priority(tick.source); + coalesced_tick_count.fetch_add(1, std::memory_order_relaxed); + if (current_priority > next_priority || + (current_priority == next_priority && + deferred_tick->issued_at >= tick.issued_at)) return; - } } + deferred_tick = tick; } void Plot::Private::arm_tick_consumer(std::weak_ptr lifetime) { - if (terminal_failure.load(std::memory_order_acquire) || - render_admission_busy.load(std::memory_order_acquire) || - !pending_tick.load(std::memory_order_acquire)) - return; + if (terminal_failure.load(std::memory_order_acquire)) return; if (tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return; aethera::schedule_task("web.plot.tick.consume", [lifetime] { const auto plot = lifetime.lock(); @@ -588,26 +596,55 @@ void Plot::Private::arm_tick_consumer(std::weak_ptr lifetime) { } void Plot::Private::release_render_admission(std::weak_ptr lifetime) { - if (!render_admission_busy.exchange(false, std::memory_order_acq_rel)) + auto expected = Render_Admission_State::rendering; + if (!render_admission.compare_exchange_strong( + expected, Render_Admission_State::ready, + std::memory_order_acq_rel, std::memory_order_acquire)) return; - if (pending_tick.load(std::memory_order_acquire)) - deferred_resume_count.fetch_add(1, std::memory_order_relaxed); + const auto pacing = frame_policy.read(); + if (pacing.render_enabled && pacing.mode == Frame_Pacing_Mode::maximum_rate) { + const auto current_consumers = consumers.load(std::memory_order_acquire); + if (!current_consumers->empty()) { + const auto now = std::chrono::steady_clock::now(); + submit_tick_request(Plot_Render_Tick{ + .issued_at = now, + .time_milliseconds = std::chrono::duration( + now - clock_origin).count(), + .source = Plot_Render_Tick_Source::maximum_rate}); + } + } arm_tick_consumer(std::move(lifetime)); } void Plot::Private::consume_tick(std::weak_ptr lifetime) { - if (!render_admission_busy.load(std::memory_order_acquire)) { - const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel); + const auto observed_generation = + tick_request_generation.load(std::memory_order_acquire); + Plot_Render_Tick requested; + while (tick_requests.try_dequeue(requested)) + keep_latest_tick(requested); + if (render_admission.load(std::memory_order_acquire) == + Render_Admission_State::ready) { + auto tick = std::exchange(deferred_tick, {}); if (tick) clock_tick(*tick); } tick_task_scheduled.store(false, std::memory_order_release); - arm_tick_consumer(std::move(lifetime)); + if (tick_request_generation.load(std::memory_order_acquire) != + observed_generation || + (render_admission.load(std::memory_order_acquire) == + Render_Admission_State::ready && deferred_tick)) + arm_tick_consumer(std::move(lifetime)); } void Plot::Private::clock_tick(const Plot_Render_Tick& tick) { if (terminal_failure.load(std::memory_order_acquire)) return; - if (tick.sequence != 0 && - !frame_policy.accept_periodic_tick(tick.time_milliseconds)) { + const auto pacing = frame_policy.read(); + const bool accepted = pacing.render_enabled && + (tick.source == Plot_Render_Tick_Source::immediate || + (tick.source == Plot_Render_Tick_Source::periodic && + pacing.mode == Frame_Pacing_Mode::fixed_rate) || + (tick.source == Plot_Render_Tick_Source::maximum_rate && + pacing.mode == Frame_Pacing_Mode::maximum_rate)); + if (!accepted) { policy_skip_count.fetch_add(1, std::memory_order_relaxed); return; } @@ -688,12 +725,12 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { const auto pacing = frame_policy.read(); if (!pacing.render_enabled || streams.consumers->empty()) return; - bool admission_expected = false; - if (!render_admission_busy.compare_exchange_strong( - admission_expected, true, std::memory_order_acq_rel, + auto admission_expected = Render_Admission_State::ready; + if (!render_admission.compare_exchange_strong( + admission_expected, Render_Admission_State::rendering, + std::memory_order_acq_rel, std::memory_order_acquire)) { - preparation_busy_count.fetch_add(1, std::memory_order_relaxed); - defer_tick(tick); + keep_latest_tick(tick); return; } @@ -737,14 +774,16 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { } if (!managed) { frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed); - defer_tick(tick); + keep_latest_tick(tick); /* * 三个槽都仍被外接消费者持有时,只保留 latest pending。这里绝不能 * 立即 arm tick consumer,否则会在没有任何槽可用期间形成 * consume -> no slot -> consume 的 Taskflow 任务风暴。真正的唤醒点 * 是 retire_completed_frame:某个 consuming 槽变回 available 后只唤醒一次。 */ - render_admission_busy.store(false, std::memory_order_release); + render_admission.store( + Render_Admission_State::frame_slots_exhausted, + std::memory_order_release); return; } managed->presentation_time = @@ -769,16 +808,28 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { const std::uint64_t sequence = next_frame_sequence++; const Frame_Identity identity{ sequence, tick.sequence == 0 ? sequence : tick.sequence}; + const auto request_source = [&] { + switch (tick.source) { + case Plot_Render_Tick_Source::periodic: + return Frame_Request_Source::periodic; + case Plot_Render_Tick_Source::immediate: + return Frame_Request_Source::immediate; + case Plot_Render_Tick_Source::maximum_rate: + return Frame_Request_Source::maximum_rate; + } + return Frame_Request_Source::unspecified; + }(); Render_Frame* logical_frame{}; if (auto* frame_2d = std::get_if>(&managed->frame)) { - (*frame_2d)->begin(identity, Frame_2D::native_pixel_format); + (*frame_2d)->begin(identity, Frame_2D::native_pixel_format, + request_source); logical_frame = frame_2d->get(); } else { auto& frame_3d = std::get>(managed->frame); frame_3d->begin(identity, pacing.video_enabled ? Frame_3D_Output::pixels : Frame_3D_Output::diagnostics, - Frame_3D::native_pixel_format); + Frame_3D::native_pixel_format, request_source); logical_frame = frame_3d.get(); } taskflow_trace_claimed = mark_taskflow_trace(*logical_frame); @@ -819,10 +870,8 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { release_render_admission(lifetime); } else { taskflow_trace_claimed = false; - frame_policy.frame_submitted(); submitted_frame_count.fetch_add(1, std::memory_order_relaxed); } - if (!result) frame_policy.frame_rejected(); return; } auto& output = *std::get>(managed->frame); @@ -832,11 +881,9 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { const auto result = scene_3d->render(&output); if (result == Render_Scene_3D::Render_Result::submitted) { taskflow_trace_claimed = false; - frame_policy.frame_submitted(); submitted_frame_count.fetch_add(1, std::memory_order_relaxed); return; } - frame_policy.frame_rejected(); scene_rejection_count.fetch_add(1, std::memory_order_relaxed); rollback_unsubmitted(); restore_taskflow_trace_claim(); @@ -1124,8 +1171,12 @@ void Plot::Private::finalize_retired_frame(Render_Frame* frame) { throw std::logic_error("retired Plot frame is not consuming"); latest_statistics_frame.store(managed, std::memory_order_release); - /* 三个消费者槽曾全部占满时,退役一个槽后继续 latest pending。 */ - arm_tick_consumer(lifetime); + /* 物理槽是唯一背压原因;归还任意槽后只解除一次耗尽状态。 */ + auto admission = Render_Admission_State::frame_slots_exhausted; + if (render_admission.compare_exchange_strong( + admission, Render_Admission_State::ready, + std::memory_order_acq_rel, std::memory_order_acquire)) + arm_tick_consumer(lifetime); if (captured_trace) { auto owner = lifetime; @@ -1175,7 +1226,10 @@ void Plot::ensure_started() { [weak](Frame_Scheduler::Tick tick) { if (const auto owner = weak.lock()) { owner->schedule_render(Plot_Render_Tick{ - tick.issued_at, tick.sequence, tick.time_milliseconds}); + .issued_at = tick.issued_at, + .sequence = tick.sequence, + .time_milliseconds = tick.time_milliseconds, + .source = Plot_Render_Tick_Source::periodic}); } }); /* @@ -1284,7 +1338,7 @@ void Plot::schedule_render(Plot_Render_Tick tick) { ensure_started(); if (d->terminal_failure.load(std::memory_order_acquire)) return; d->received_tick_count.fetch_add(1, std::memory_order_relaxed); - d->defer_tick(tick); + d->submit_tick_request(std::move(tick)); d->arm_tick_consumer(weak_from_this()); } @@ -1294,7 +1348,10 @@ void Plot::render_once() { const auto now = std::chrono::steady_clock::now(); const auto elapsed = now - d->clock_origin; schedule_render(Plot_Render_Tick{ - now, 0, std::chrono::duration(elapsed).count()}); + .issued_at = now, + .time_milliseconds = + std::chrono::duration(elapsed).count(), + .source = Plot_Render_Tick_Source::immediate}); } void Plot::submit_input(Plot_Input_Event event) { @@ -1328,7 +1385,7 @@ nlohmann::json Plot::write_prop(std::string_view component, ensure_started(); if (component != "frame-analysis") return d->view->write_prop(component, key, value); - auto result = d->frame_policy.write_prop(key, value); + auto result = write_frame_policy_prop(d->frame_policy, key, value); if (result.value("success", false)) d->refresh_schedule(); return result; } @@ -1395,6 +1452,16 @@ nlohmann::json Plot::diagnostics() const { const auto pacing = d->frame_policy.read(); const auto stream = d->stream_snapshot(); + const auto admission = d->render_admission.load(std::memory_order_acquire); + const auto admission_name = [&] { + switch (admission) { + case Private::Render_Admission_State::ready: return "ready"; + case Private::Render_Admission_State::rendering: return "rendering"; + case Private::Render_Admission_State::frame_slots_exhausted: + return "frame_slots_exhausted"; + } + return "unknown"; + }(); nlohmann::json supported_formats = nlohmann::json::array(); if (is_3d) { for (const auto format : Frame_3D::supported_pixel_formats) @@ -1436,9 +1503,7 @@ nlohmann::json Plot::diagnostics() const { {"received_ticks", d->received_tick_count.load(std::memory_order_relaxed)}, {"coalesced_ticks", d->coalesced_tick_count.load(std::memory_order_relaxed)}, {"policy_skips", d->policy_skip_count.load(std::memory_order_relaxed)}, - {"preparation_busy", d->preparation_busy_count.load(std::memory_order_relaxed)}, - {"deferred_resumes", d->deferred_resume_count.load(std::memory_order_relaxed)}, - {"render_admission_busy", d->render_admission_busy.load(std::memory_order_relaxed)}, + {"render_admission", admission_name}, {"frame_slot_busy", d->frame_slot_busy_count.load(std::memory_order_relaxed)}, {"scene_rejections", d->scene_rejection_count.load(std::memory_order_relaxed)}, {"submitted_frames", d->submitted_frame_count.load(std::memory_order_relaxed)}}}, diff --git a/web_server/src/Plot.hpp b/web_server/src/Plot.hpp index 2969584..7e31def 100644 --- a/web_server/src/Plot.hpp +++ b/web_server/src/Plot.hpp @@ -27,12 +27,19 @@ struct Plot_Input_Event { std::uint32_t native_key{}; /* 浏览器原生按键码。 */ bool auto_repeat{}; /* 是否为系统重复按键。 */ }; +enum struct Plot_Render_Tick_Source : std::uint8_t { + periodic, + immediate, + maximum_rate +}; + struct Plot_Render_Tick { std::chrono::steady_clock::time_point issued_at{}; /* 页面帧时钟发布本 tick 的单调时刻;手动帧在提交时填写。 */ std::uint64_t sequence{}; /* Kernel 全局 Frame_Scheduler 时间轴上的关联序号。 */ double time_milliseconds{}; /* 页面级单调时间线,所有图共享同一个动画时刻。 */ std::uint32_t width{320}; /* 当前图在媒体图集中的固定像素宽度。 */ std::uint32_t height{192}; /* 当前图在媒体图集中的固定像素高度。 */ + Plot_Render_Tick_Source source{Plot_Render_Tick_Source::immediate}; /* 本次请求来自周期时钟、手动操作或最大吞吐自驱动。 */ }; enum struct Plot_Pixel_Layout : std::uint8_t { bgra8, diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index b3fde7b..927a20c 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -85,6 +85,7 @@ type Taskflow_Execution_Trace = {native_id: string; node_id: string; worker_id: cpu_time_coarse: boolean; observer_entry_ms: number; observer_exit_ms: number; observer_entry_cpu_ms: number; observer_exit_cpu_ms: number; queue_wait_ms: number}; type Taskflow_Frame_Trace = {sequence: number; correlation_id: number; created_time_unix_ns: number; worker_count: number; + request_source?: "unspecified" | "periodic" | "immediate" | "maximum_rate"; markers?: Record; measurements?: Record; graphs: Taskflow_Graph_Trace[]; executions: Taskflow_Execution_Trace[]}; type Taskflow_Frame_Response = {protocol: "aethera.taskflow.frames"; version: 1; requested: number; remaining: number; @@ -1839,6 +1840,7 @@ function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_ await copy_text(JSON.stringify({ sequence: frame?.sequence, correlation_id: frame?.correlation_id, + request_source: frame?.request_source, markers: frame?.markers ?? {}, measurements: frame?.measurements ?? {}, graph, @@ -1972,7 +1974,8 @@ function Taskflow_Timeline({graph, executions, frame, components, gallery_state, groups.push({id, title:
{title}{detail}
, node_name: id, type: "lifecycle"}); if (frame) { - add_lifecycle_group("lifecycle-plot", "Plot 调度与更新", "tick → 获得帧槽 → view.update"); + add_lifecycle_group("lifecycle-plot", "Plot 调度与更新", + `${frame.request_source ?? "unspecified"} → 获得帧槽 → view.update`); const update_started = marker("plot_update_started") ?? 0; const tick_queue = measurements.plot_tick_queue_ns ?? 0; add_item("lifecycle-tick-queue", "lifecycle-plot", "lifecycle", @@ -2112,6 +2115,7 @@ function Taskflow_Timeline({graph, executions, frame, components, gallery_state, await copy_text(JSON.stringify({ sequence: frame?.sequence, correlation_id: frame?.correlation_id, + request_source: frame?.request_source, markers: frame?.markers ?? {}, measurements: frame?.measurements ?? {}, graph,