diff --git a/kernel/src/kernel/Frame_Pacing_Policy.cpp b/kernel/src/kernel/Frame_Pacing_Policy.cpp index 8688e8e..56515a9 100644 --- a/kernel/src/kernel/Frame_Pacing_Policy.cpp +++ b/kernel/src/kernel/Frame_Pacing_Policy.cpp @@ -59,31 +59,23 @@ bool Frame_Pacing_Policy::accept_periodic_tick(double time_milliseconds) { return false; if (skip_next_periodic_.exchange(false, std::memory_order_acq_rel)) return false; - return !in_flight_.load(std::memory_order_acquire); + 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); - if (in_flight_.load(std::memory_order_acquire)) { - urgent_pending_.store(true, std::memory_order_release); - return false; - } return true; } -void Frame_Pacing_Policy::frame_submitted() { - in_flight_.store(true, std::memory_order_release); -} +void Frame_Pacing_Policy::frame_submitted() {} bool Frame_Pacing_Policy::frame_completed() { - in_flight_.store(false, std::memory_order_release); - return urgent_pending_.exchange(false, std::memory_order_acq_rel); + return false; } -void Frame_Pacing_Policy::frame_rejected() { - in_flight_.store(false, std::memory_order_release); -} +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 2c30f28..fc3731a 100644 --- a/kernel/src/kernel/Frame_Pacing_Policy.hpp +++ b/kernel/src/kernel/Frame_Pacing_Policy.hpp @@ -15,7 +15,7 @@ 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{100.0}; + double fixed_rate_fps{30.0}; double maximum_rate_fps{100.0}; }; @@ -36,8 +36,8 @@ struct Frame_Pacing_Policy final { [[nodiscard]] double scheduled_rate_fps() const; /* - * 周期 tick 是否允许生成新帧。immediate 会提前生成一帧并消费下一次周期机会。 - * in_flight 只表达 Scene 是否已有一帧尚未完成,不改变 Scene::render(Frame*) 抽象。 + * 周期 tick 只判断策略/周期语义,不再承担 Scene in-flight 门控。 + * 实际渲染准入由调用方在 Scene::advance 前控制;忙时保留 latest pending tick。 */ [[nodiscard]] bool accept_periodic_tick(double time_milliseconds); [[nodiscard]] bool request_immediate(); @@ -49,9 +49,7 @@ private: template void update(Edit&& edit); std::atomic> properties_; - std::atomic_bool in_flight_{}; std::atomic_bool skip_next_periodic_{}; - std::atomic_bool urgent_pending_{}; }; } diff --git a/kernel/src/kernel/Frame_Scheduler.cpp b/kernel/src/kernel/Frame_Scheduler.cpp index 8dd71eb..f077e64 100644 --- a/kernel/src/kernel/Frame_Scheduler.cpp +++ b/kernel/src/kernel/Frame_Scheduler.cpp @@ -137,7 +137,17 @@ struct Frame_Scheduler::Private { if (!state->alive.load(std::memory_order_acquire)) break; state->mode = Timer::State::Mode::periodic; state->period = command.value; - state->next_deadline = now + state->period; + /* 同频 Plot 不要永久锁在同一 deadline 上。稳定 timer id 给 + * 每个周期分配固定 phase;时间值仍来自同一全局 origin。 */ + { + constexpr std::uint64_t phase_slots = 8; + const auto phase_slot = (state->id - 1U) % phase_slots; + const auto phase = Nanoseconds{ + state->period.count() * + static_cast(phase_slot) / + static_cast(phase_slots)}; + state->next_deadline = now + state->period + phase; + } schedule_state(*state, now); break; case Command_Type::once: @@ -169,13 +179,14 @@ struct Frame_Scheduler::Private { void fire(Timer::State& state) noexcept { if (!state.alive.load(std::memory_order_acquire)) return; - const auto now = Scheduler_Clock::now(); + const auto issued_at = Scheduler_Clock::now(); const ::Tick sequence = wheel.now(); try { state.handler(Frame_Scheduler::Tick{ - now, + issued_at, static_cast(sequence), - std::chrono::duration(now - origin).count()}); + std::chrono::duration( + issued_at - origin).count()}); } catch (...) { /* Timer callback 是控制面入口,不能因为业务异常终止全局时钟。 */ @@ -189,12 +200,16 @@ struct Frame_Scheduler::Private { state.period <= Nanoseconds::zero()) return; + /* callback 自身也可能消耗控制线程时间;重取 now,避免用触发前时间 + * 重排 deadline 导致已经过期的周期再次被排入 timer wheel。 */ + const auto reschedule_now = Scheduler_Clock::now(); state.next_deadline += state.period; - if (state.next_deadline <= now) { - const auto missed = (now - state.next_deadline) / state.period + 1; + if (state.next_deadline <= reschedule_now) { + const auto missed = + (reschedule_now - state.next_deadline) / state.period + 1; state.next_deadline += state.period * missed; } - schedule_state(state, now); + schedule_state(state, reschedule_now); } void advance_to(Scheduler_Clock::time_point now) { diff --git a/kernel/src/kernel/render_common.cpp b/kernel/src/kernel/render_common.cpp index 89f1c0b..4516f80 100644 --- a/kernel/src/kernel/render_common.cpp +++ b/kernel/src/kernel/render_common.cpp @@ -156,8 +156,8 @@ private: std::vector worker_cpu_starts; std::unique_ptr worker_statistics; std::size_t worker_statistics_count{}; - std::atomic trace_frame{}; - std::shared_mutex trace_mutex{}; /* 仅按需捕获时保护 Frame* 获取与关闭。 */ + std::unordered_map trace_tasks{}; + std::shared_mutex trace_mutex{}; /* 按原生节点身份把并发 Taskflow 归属到各自 Frame。 */ std::uint64_t worker_occupation_limit_ns{}; /* 单节点连续非 CPU 等待 Worker 的上限。 */ Task_Overrun_Action worker_overrun_action{}; /* 节点超过占用上限后的处置策略。 */ std::atomic_bool watchdog_stopping{}; /* Watchdog 生命周期停止标志。 */ @@ -433,22 +433,20 @@ public: record.native_id = static_cast(task.hash_value()); record.type = task.type(); worker_starts.push_back(std::move(record)); - auto* frame = trace_frame.load(std::memory_order_acquire); + Render_Frame* frame{}; /* - * 帧租约从 on_entry 持续到对应 on_exit。只在退出时登记写入者会留下 - * “捕获已关闭、物理帧已复用、迟到 on_exit 仍访问旧 Frame*”的窗口。 + * 每次被追踪的 Taskflow run 在提交前登记 native node -> Frame。 + * 因而 Render(N+1) 与外接 H264(N) 可以同时被 Observer 正确归属, + * 不再使用“全局唯一当前 Frame”的假设。帧租约从 on_entry 持续到 + * on_exit,避免 graph 退役后迟到的 Observer 写入访问已复用 Frame。 */ - if (frame) { + { std::shared_lock trace_guard(trace_mutex); - frame = trace_frame.load(std::memory_order_acquire); - if (frame && detail::Taskflow_Frame_Access::acquire_writer(*frame)) { - if (!detail::Taskflow_Frame_Access::contains_task( - *frame, static_cast(task.hash_value()))) { - detail::Taskflow_Frame_Access::release_writer(*frame); - frame = nullptr; - } - } - else frame = nullptr; + const auto found = trace_tasks.find( + static_cast(task.hash_value())); + if (found != trace_tasks.end() && + detail::Taskflow_Frame_Access::acquire_writer(*found->second)) + frame = found->second; } worker_starts.back().frame = frame; const auto active_depth = worker_state.active_depth.fetch_add( @@ -645,17 +643,44 @@ public: update_max(worker_state.last_task_time_ns, clock_ns(completed)); } bool begin_trace(Render_Frame& frame, std::size_t workers) { - std::unique_lock guard(trace_mutex); - const auto current = trace_frame.load(std::memory_order_acquire); - if (current) return current == &frame; + /* Frame 自己保存捕获缓冲;多个 Frame 可以同时处于捕获窗口。 */ detail::Taskflow_Frame_Access::begin_capture(frame, workers); - trace_frame.store(&frame, std::memory_order_release); return true; } - void finish_trace(Render_Frame& frame) noexcept { + void activate_graph(Render_Frame& frame, Task_Graph& graph) { + if (!frame.taskflow_trace_requested()) return; + const auto nodes = detail::Task_Graph_Access::nodes(graph); std::unique_lock guard(trace_mutex); - if (trace_frame.load(std::memory_order_acquire) != &frame) return; - trace_frame.store(nullptr, std::memory_order_release); + for (const auto& node : nodes) { + const auto [found, inserted] = trace_tasks.emplace( + node.native_id, &frame); + if (!inserted && found->second != &frame) + throw std::logic_error( + "concurrent traced runs share the same Taskflow node"); + } + } + void deactivate_graph(Render_Frame& frame, Task_Graph& graph) noexcept { + if (!frame.taskflow_trace_requested()) return; + try { + const auto nodes = detail::Task_Graph_Access::nodes(graph); + std::unique_lock guard(trace_mutex); + for (const auto& node : nodes) { + const auto found = trace_tasks.find(node.native_id); + if (found != trace_tasks.end() && found->second == &frame) + trace_tasks.erase(found); + } + } + catch (...) {} + } + void finish_trace(Render_Frame& frame) noexcept { + try { + std::unique_lock guard(trace_mutex); + for (auto it = trace_tasks.begin(); it != trace_tasks.end();) { + if (it->second == &frame) it = trace_tasks.erase(it); + else ++it; + } + } + catch (...) {} detail::Taskflow_Frame_Access::finish_capture(frame); } void write_state(Task_Runtime_State& state, std::size_t workers, std::size_t active_topologies) const { @@ -929,12 +954,15 @@ public: std::string_view stage) { const auto token = detail::Taskflow_Frame_Access::begin_graph( frame, taskflow, stage); + observer->activate_graph(frame, taskflow); try { const auto elapsed = run(taskflow); + observer->deactivate_graph(frame, taskflow); detail::Taskflow_Frame_Access::finish_graph(token); return elapsed; } catch (...) { + observer->deactivate_graph(frame, taskflow); detail::Taskflow_Frame_Access::finish_graph(token); throw; } @@ -943,10 +971,20 @@ public: std::function completion) { const auto token = detail::Taskflow_Frame_Access::begin_graph( frame, taskflow, stage); - run(taskflow, [token, completion = std::move(completion)]() mutable { + try { + observer->activate_graph(frame, taskflow); + run(taskflow, [this, &taskflow, &frame, token, + completion = std::move(completion)]() mutable { + observer->deactivate_graph(frame, taskflow); + detail::Taskflow_Frame_Access::finish_graph(token); + completion(); + }); + } + catch (...) { + observer->deactivate_graph(frame, taskflow); detail::Taskflow_Frame_Access::finish_graph(token); - completion(); - }); + throw; + } } void schedule(std::string name, std::function task) { if (!task) throw std::invalid_argument("Taskflow scheduled task is empty"); diff --git a/render_2D/render_2D/scene/Render_Scene_2D.cpp b/render_2D/render_2D/scene/Render_Scene_2D.cpp index 7fddf19..2d50c71 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.cpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.cpp @@ -9,6 +9,10 @@ Render_Scene_2D::render(Frame_2D* frame) { return static_cast(*d).dispatch->render(this, frame); } void Render_Scene_2D::set_frame_callback(Frame_Callback callback) { static_cast(*d).dispatch->set_frame_callback(this, std::move(callback)); } +void Render_Scene_2D::set_frame_retired_callback(Frame_Callback callback) { + static_cast(*d).dispatch->set_frame_retired_callback( + this, std::move(callback)); +} Task_Graph& Render_Scene_2D::completion_taskflow() { return static_cast(*d).completion_graph; } diff --git a/render_2D/render_2D/scene/Render_Scene_2D.hpp b/render_2D/render_2D/scene/Render_Scene_2D.hpp index b436ad9..814c8e0 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.hpp +++ b/render_2D/render_2D/scene/Render_Scene_2D.hpp @@ -43,13 +43,16 @@ struct Render_Scene_2D : Def; - /* 向调用方拥有的帧合成一次;必须先安装回调。回调返回前的并发请求返回 frame_in_flight。 */ + /* 向调用方拥有的帧合成一次;Scene::advance 到完成像素发布保持不可重入。 */ [[nodiscard]] std::expected render(Frame_2D* frame); /* - * 安装完成帧回调;回调收到对应 render(frame) 的对象,返回时 Scene 才释放下一帧准入。 - * Scene 和调用方拥有的 Frame 必须存活到回调返回;析构不会用同步等待隐藏错误的生命周期。 + * 安装完成帧回调;回调在 Scene 像素/发布阶段结束并释放下一帧准入后执行。 + * 因而回调可以继续执行外接编码/发送 Taskflow,同时下一帧 Scene 已可开始。 + * 调用方拥有的 Frame 必须存活到回调返回。 */ void set_frame_callback(Frame_Callback callback); + /* callback/frame_ready/trace 完全收尾后通知调用方归还物理 Frame 槽。 */ + void set_frame_retired_callback(Frame_Callback callback); /* * 返回最终像素完成后、帧回调前执行的直接 Taskflow。 * 只能在 Scene 没有运行时修改该图;禁止在执行期间 emplace/erase/clear。 diff --git a/render_2D/render_2D/scene/Render_Scene_2D.ipp b/render_2D/render_2D/scene/Render_Scene_2D.ipp index c3e6740..b510cce 100644 --- a/render_2D/render_2D/scene/Render_Scene_2D.ipp +++ b/render_2D/render_2D/scene/Render_Scene_2D.ipp @@ -67,19 +67,19 @@ struct Render_Scene_2D::Private : Prev_Private { struct Dispatch { void (*reset_statistics)(Root*); Render_Run render; /* 向外部帧执行最终 Scene 并合成颜色层。 */ - Callback_Run set_frame_callback; /* 安装最终完成帧回调。 */ + Callback_Run set_frame_callback; /* 安装像素完成后的外接消费回调。 */ + Callback_Run set_frame_retired_callback; /* trace/frame_ready 完成后的物理帧退役回调。 */ Active_Run set_active; /* 修改最终 Scene 的视图活动状态。 */ }; const Dispatch* dispatch{}; /* Builder 绑定最终 Scene 类型后的静态分派表。 */ std::mutex render_mutex{}; /* 只保护完成回调和单帧准入。 */ - bool frame_in_flight{}; /* render 准入到完成回调返回的唯一状态源。 */ - Frame_Callback frame_callback{}; /* 合成完成后的唯一像素发布出口。 */ + bool frame_in_flight{}; /* Scene::advance 到像素发布完成的唯一不可重入准入状态。 */ + Frame_Callback frame_callback{}; /* Plot publish 后的外接消费出口。 */ + Frame_Callback frame_retired_callback{}; /* 全链路 trace 完成后的物理帧归还出口。 */ Frame_Statistics_Accumulator frame_statistics{}; /* Scene 内部增量计算;State 只发布定长统计结果。 */ Task_Graph completion_graph{"render_2d.completion"}; /* 最终像素完成后、发布回调前执行的外部续写图。 */ std::unique_ptr paint_taskflow{}; /* 仅由二维 Paint 图构建的执行图。 */ std::unique_ptr frame_taskflow{}; - Frame_Callback active_callback{}; - bool active_trace{}; std::vector active_prepare_executions{}; std::chrono::steady_clock::time_point active_prepare_started{}; ~Private(); @@ -372,6 +372,7 @@ template std::expected Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { Frame_Callback callback; + Frame_Callback retired_callback; { std::lock_guard lock(render_mutex); if (!frame_callback) @@ -380,6 +381,7 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { return std::unexpected(Render_Result::frame_in_flight); frame_in_flight = true; callback = frame_callback; + retired_callback = frame_retired_callback; } const auto release_admission = [this] { std::lock_guard lock(render_mutex); @@ -407,41 +409,52 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { } ensure_frame_taskflow(object); active_frame = frame; - active_callback = callback; } catch (...) { release_admission(); throw; } frame->mark(Frame_Trace_Marker::scene_render_requested); - active_trace = aethera::detail::begin_taskflow_trace(*frame); + const bool trace_started = aethera::detail::begin_taskflow_trace(*frame); try { - auto completion = [this, object, frame] { - struct Release { - Private* data; - ~Release() { - data->active_frame = nullptr; - data->frame_target = nullptr; - data->active_callback = {}; - data->active_prepare_executions.clear(); - std::lock_guard lock(data->render_mutex); - data->frame_in_flight = false; - } - } release{this}; - if (active_trace) { - aethera::detail::finish_taskflow_trace(*frame); - active_trace = false; - } - frame->mark(Frame_Trace_Marker::callback_started); - active_callback(frame); - frame->mark(Frame_Trace_Marker::callback_finished); - frame->mark(Frame_Trace_Marker::frame_ready); + auto completion = [this, object, frame, callback = std::move(callback), + retired_callback = std::move(retired_callback), + trace_started]() mutable { + /* + * render_2d.frame 的原生 Taskflow 已经完成,所有 active_* 只属于 + * Scene::advance -> pixel/publish 的不可重入区。先发布 Scene 自身 + * 渲染统计并释放准入,再进入调用方 callback;callback 可以继续 + * corun 外接 H264/WebRTC DAG,而下一帧已经允许开始。 + */ auto& private_data = static_cast(*this); auto& scene_state = static_cast(*private_data.state.pending); scene_state.frame_statistics = frame_statistics.submit( *frame, Frame_Dimension::two_dimensional); private_data.state.advance(); object->template notify_state(); + + active_frame = nullptr; + frame_target = nullptr; + active_prepare_executions.clear(); + { + std::lock_guard lock(render_mutex); + frame_in_flight = false; + } + + try { + frame->mark(Frame_Trace_Marker::callback_started); + callback(frame); + frame->mark(Frame_Trace_Marker::callback_finished); + frame->mark(Frame_Trace_Marker::frame_ready); + } + catch (...) { + if (trace_started) + aethera::detail::finish_taskflow_trace(*frame); + throw; + } + if (trace_started) + aethera::detail::finish_taskflow_trace(*frame); + if (retired_callback) retired_callback(frame); }; if (frame->taskflow_trace_requested()) aethera::detail::run_taskflow( @@ -450,13 +463,10 @@ Render_Scene_2D::Private::render(Object* object, Frame_2D* frame) { aethera::detail::run_taskflow(*frame_taskflow, std::move(completion)); } catch (...) { - if (active_trace) { + if (trace_started) aethera::detail::finish_taskflow_trace(*frame); - active_trace = false; - } active_frame = nullptr; frame_target = nullptr; - active_callback = {}; active_prepare_executions.clear(); std::lock_guard lock(render_mutex); frame_in_flight = false; @@ -535,6 +545,12 @@ const Render_Scene_2D::Private::Dispatch& Render_Scene_2D::Private::dispatch_for std::lock_guard lock(data.render_mutex); data.frame_callback = std::move(callback); }, + [](Root* root, Frame_Callback callback) { + auto* object = static_cast(root); + auto& data = static_cast(*object->d); + std::lock_guard lock(data.render_mutex); + data.frame_retired_callback = std::move(callback); + }, [](Root* root, bool active) { static_cast(root)->template set<&Prop::view_active>(active); } }; return value; diff --git a/web_server/src/Gallery_Video_Stream.cpp b/web_server/src/Gallery_Video_Stream.cpp index c5ea78a..666de52 100644 --- a/web_server/src/Gallery_Video_Stream.cpp +++ b/web_server/src/Gallery_Video_Stream.cpp @@ -1,7 +1,7 @@ #include "Gallery_Video_Stream.hpp" #include #include -#include "detail/Gallery_Frame_Atlas.hpp" +#include "frame_sampling/Frame_Sampler.hpp" #include #include #include @@ -18,7 +18,7 @@ namespace { constexpr std::uint32_t tile_width{720}; constexpr std::uint32_t tile_height{420}; constexpr std::uint32_t atlas_columns{4}; -constexpr double gallery_frame_rate{100.0}; +constexpr double gallery_frame_rate{30.0}; constexpr auto metric_interval{std::chrono::seconds(1)}; nlohmann::json statistic_json(const Statistic_State& value) { return { @@ -43,7 +43,8 @@ std::string exception_description(const std::exception_ptr& failure) { return "empty gallery video failure"; } } -Gallery_Video_Stream::Private::Private() : encoder(gallery_frame_rate) {} +Gallery_Video_Stream::Private::Private() + : ffmpeg_transport(gallery_frame_rate) {} Gallery_Video_Stream::Private::~Private() = default; void Gallery_Video_Stream::Private::initialize( std::vector plots) { @@ -56,17 +57,20 @@ void Gallery_Video_Stream::Private::initialize( source_ids.push_back(plot.id); sources.push_back(Source{std::move(plot)}); } - atlas = std::make_unique( - tile_width, tile_height, atlas_columns, std::move(source_ids)); + 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); metric_completion_starts.resize(sources.size()); metric_rendered_starts.resize(sources.size()); } -bool Gallery_Video_Stream::Private::consumer_accepts_video() const noexcept { +bool Gallery_Video_Stream::Private::consumer_accepts( + Gallery_Transport_Mode transport) const noexcept { const auto current = consumers.load(std::memory_order_acquire); for (const auto& consumer : *current) { - if (!consumer.video_readiness) continue; + if (consumer.transport != transport || !consumer.readiness) continue; try { - if ((*consumer.video_readiness)()) return true; + if ((*consumer.readiness)()) return true; } catch (...) {} } @@ -90,16 +94,19 @@ bool Gallery_Video_Stream::Private::remove_consumer( } } void Gallery_Video_Stream::Private::publish( - std::optional video, + Gallery_Transport_Mode transport, + std::shared_ptr video, + std::shared_ptr pixels, std::string notification) noexcept { - if (!video && notification.empty()) return; + if (!video && !pixels && notification.empty()) return; try { const auto current = consumers.load(std::memory_order_acquire); std::vector failed_consumers; for (const auto& consumer : *current) { - if (!consumer.handler) continue; + if (consumer.transport != transport || !consumer.handler) continue; try { - (*consumer.handler)(Gallery_Stream_Frame{video, notification}); + (*consumer.handler)(Gallery_Stream_Frame{ + video, pixels, notification}); } catch (...) { failed_consumers.push_back(consumer.id); @@ -115,12 +122,14 @@ void Gallery_Video_Stream::Private::fail( if (failed.exchange(true, std::memory_order_acq_rel)) return; try { terminal_failure = exception_description(failure); - publish(std::nullopt, nlohmann::json{ + const auto notification = nlohmann::json{ {"kind", "gallery_error"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"message", *terminal_failure} - }.dump()); + }.dump(); + publish(Gallery_Transport_Mode::ffmpeg, {}, {}, notification); + publish(Gallery_Transport_Mode::websocket_pixels, {}, {}, notification); } catch (...) {} } @@ -129,24 +138,26 @@ void Gallery_Video_Stream::Private::accept_frame( if (stopping.load(std::memory_order_acquire) || failed.load(std::memory_order_acquire) || !frame) return; - static_cast(atlas->accept_frame(slot, frame->pixels)); - if (slot == encode_source_slot && frame->pixels) { + 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); - active_encode_tick.sequence = frame->pixels->correlation_id; - active_encode_tick.time_milliseconds = - static_cast(frame->pixels->presentation_time.count()) / - 1'000.0; } } void Gallery_Video_Stream::Private::update_metrics( const Plot_Render_Tick& tick, - const detail::Gallery_Atlas_Composition& composition, + const frame_sampling::Sampled_Gallery_Frame& sample, std::size_t encoded_bytes, - Video_Encoder_Backend encoder_backend) { + std::size_t pixel_bytes, + std::optional encoder_backend) { + const auto& composition = sample.composition; + const auto sampler_state = sampler->state(); const auto now = std::chrono::steady_clock::now(); if (metric_started.time_since_epoch().count() == 0) { metric_started = now; 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); for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) { metric_completion_starts[slot] = @@ -162,7 +173,7 @@ void Gallery_Video_Stream::Private::update_metrics( const auto clock_total = delivered_clock_ticks.load(std::memory_order_relaxed); nlohmann::json source_metrics = nlohmann::json::object(); - const auto description = atlas->describe(); + const auto description = sampler->describe(); for (std::size_t slot = 0; slot < composition.sources.size(); ++slot) { const auto& progress = composition.sources[slot]; const auto completed = progress.completion_count - @@ -199,71 +210,110 @@ void Gallery_Video_Stream::Private::update_metrics( metric_rendered_starts[slot] = progress.rendered_frame_count; } const auto& compose = compose_ms.state(); + const auto& sample_delay = sample_delay_ms.state(); const auto& encode = encode_ms.state(); + const auto& pixel_pack = pixel_pack_ms.state(); const auto& publish_time = publish_ms.state(); auto output = nlohmann::json{ {"kind", "gallery_metrics"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"clock_sequence", tick.sequence}, + {"sampled_frame_count", sampler_state.sampled_frames}, {"encoded_frame_count", encoded_frame_count}, - {"encoder_backend", video_encoder_backend_name(encoder_backend)}, + {"pixel_frame_count", pixel_frame_count}, + {"encoder_backend", encoder_backend + ? video_encoder_backend_name(*encoder_backend) + : std::string_view{"inactive"}}, {"profile_level_id", h264_profile_level_id()}, {"target_frame_rate_fps", gallery_frame_rate}, { "clock_delivery_rate_fps", static_cast(clock_total - metric_clock_start) / seconds }, + { + "sampled_frame_rate_fps", + static_cast(sampler_state.sampled_frames - + metric_sample_start) / seconds + }, { "encoded_frame_rate_fps", static_cast(encoded_frame_count - metric_encoded_start) / seconds }, + { + "pixel_frame_rate_fps", + static_cast(pixel_frame_count - metric_pixel_start) / seconds + }, {"compose_average_ms", compose.average}, {"compose_p95_ms", compose.p95}, + {"sample_delay_average_ms", sample_delay.average}, + {"sample_delay_p95_ms", sample_delay.p95}, {"encode_average_ms", encode.average}, {"encode_p95_ms", encode.p95}, + {"pixel_pack_average_ms", pixel_pack.average}, + {"pixel_pack_p95_ms", pixel_pack.p95}, {"publish_average_ms", publish_time.average}, {"publish_p95_ms", publish_time.p95}, { "statistics", { {"compose_ms", statistic_json(compose)}, + {"sample_delay_ms", statistic_json(sample_delay)}, {"encode_ms", statistic_json(encode)}, + {"pixel_pack_ms", statistic_json(pixel_pack)}, {"publish_ms", statistic_json(publish_time)} } }, {"encoded_bytes", encoded_bytes}, + {"pixel_bytes", pixel_bytes}, {"fresh_tiles", composition.fresh_tile_count}, {"missing_tiles", composition.missing_tile_count}, {"rejected_frames", composition.rejected_frame_count}, { - "skipped_encode_ticks", - skipped_encode_ticks.load(std::memory_order_relaxed) + "skipped_sample_ticks", + skipped_sample_ticks.load(std::memory_order_relaxed) }, + {"sampler", { + {"accepted_source_frames", sampler_state.accepted_source_frames}, + {"rejected_source_frames", sampler_state.rejected_source_frames}, + {"early_ticks", sampler_state.early_ticks}, + {"missed_periods", sampler_state.missed_periods}}}, {"sources", std::move(source_metrics)} }; metric_started = now; metric_encoded_start = encoded_frame_count; + metric_pixel_start = pixel_frame_count; + metric_sample_start = sampler_state.sampled_frames; metric_clock_start = clock_total; object->template update_state<&State::diagnostics>(std::move(output)); object->template publish_state(); } -void Gallery_Video_Stream::Private::compose_media_frame() { - active_composition.reset(); +void Gallery_Video_Stream::Private::sample_media_frame() { + active_sample.reset(); active_video.reset(); + active_pixels.reset(); if (stopping.load(std::memory_order_acquire) || failed.load(std::memory_order_acquire) || - !consumer_accepts_video()) { - skipped_encode_ticks.fetch_add(1, std::memory_order_relaxed); + (!consumer_accepts(Gallery_Transport_Mode::ffmpeg) && + !consumer_accepts(Gallery_Transport_Mode::websocket_pixels))) { + skipped_sample_ticks.fetch_add(1, std::memory_order_relaxed); return; } const auto started = std::chrono::steady_clock::now(); - active_composition = atlas->compose(); + active_sample = sampler->sample(); + if (!active_sample) { + skipped_sample_ticks.fetch_add(1, std::memory_order_relaxed); + return; + } static_cast(compose_ms.submit( std::chrono::duration( std::chrono::steady_clock::now() - started).count())); + static_cast(sample_delay_ms.submit( + active_sample->deadline_delay_ms)); } void Gallery_Video_Stream::Private::encode_media_frame() { - if (!active_composition) return; + if (!active_sample || + !consumer_accepts(Gallery_Transport_Mode::ffmpeg)) + return; /* * 直接在 FFmpeg.H264.encode Taskflow 节点中调用 FFmpeg。这样实际 * avcodec_send_frame/avcodec_receive_packet(以及硬件后端等待)全部落在 @@ -272,38 +322,56 @@ void Gallery_Video_Stream::Private::encode_media_frame() { * H264_Encoder 因而仍只有一个执行位置,不需要额外 mutex 或专用线程。 */ const auto started = std::chrono::steady_clock::now(); - if (key_frame_requested.exchange(false, - std::memory_order_acq_rel)) - encoder.request_key_frame(); - active_video = encoder.encode( - active_composition->pixels, active_composition->width, - active_composition->height, - active_composition->layout == Plot_Pixel_Layout::bgra8 - ? Video_Pixel_Layout::bgra - : Video_Pixel_Layout::rgba, - active_encode_tick.sequence, - std::chrono::microseconds{ - static_cast( - std::llround(active_encode_tick.time_milliseconds * - 1'000.0)) - }); + if (auto encoded = ffmpeg_transport.encode(*active_sample)) + active_video = std::make_shared( + std::move(*encoded)); static_cast(encode_ms.submit( std::chrono::duration( std::chrono::steady_clock::now() - started).count())); } -void Gallery_Video_Stream::Private::publish_media_frame() { - if (!active_video || !active_composition) return; - ++encoded_frame_count; - const auto encoded_bytes = active_video->annex_b.size(); - const auto encoder_backend = active_video->backend; - update_metrics(active_encode_tick, *active_composition, - encoded_bytes, encoder_backend); +void Gallery_Video_Stream::Private::pack_pixel_frame() { + if (!active_sample || + !consumer_accepts(Gallery_Transport_Mode::websocket_pixels)) + return; const auto started = std::chrono::steady_clock::now(); - publish(std::move(active_video), {}); + active_pixels = std::make_shared( + frame_sampling::websocket_pixel::pack_websocket_pixel_frame( + *active_sample)); + static_cast(pixel_pack_ms.submit( + std::chrono::duration( + std::chrono::steady_clock::now() - started).count())); +} +void Gallery_Video_Stream::Private::publish_ffmpeg_frame() { + if (!active_video) return; + ++encoded_frame_count; + const auto started = std::chrono::steady_clock::now(); + publish(Gallery_Transport_Mode::ffmpeg, active_video, {}, {}); static_cast(publish_ms.submit( std::chrono::duration( std::chrono::steady_clock::now() - started).count())); - active_composition.reset(); +} +void Gallery_Video_Stream::Private::publish_pixel_frame() { + if (!active_pixels) return; + ++pixel_frame_count; + const auto started = std::chrono::steady_clock::now(); + publish(Gallery_Transport_Mode::websocket_pixels, {}, active_pixels, {}); + static_cast(publish_ms.submit( + std::chrono::duration( + std::chrono::steady_clock::now() - started).count())); +} +void Gallery_Video_Stream::Private::complete_media_frame() { + if (!active_sample) return; + const auto encoded_bytes = active_video ? active_video->annex_b.size() : 0U; + const auto pixel_bytes = active_pixels ? active_pixels->bytes.size() : 0U; + const auto encoder_backend = active_video + ? std::optional{active_video->backend} + : std::nullopt; + update_metrics(active_sample->tick, *active_sample, encoded_bytes, + pixel_bytes, encoder_backend); + active_sample.reset(); + active_video.reset(); + active_pixels.reset(); } std::shared_ptr Gallery_Video_Stream::create( std::vector plots) { @@ -323,6 +391,87 @@ Gallery_Video_Stream::~Gallery_Video_Stream() { void Gallery_Video_Stream::bind_plots() { auto& data = static_cast(*d); const auto weak = weak_from_this(); + if (data.sources.empty()) + 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 的初始化窗口。 + */ + auto media = std::make_unique("gallery.video.frame"); + auto sample = media->add("gallery.sample.capture", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.sample_media_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + sample.describe("owner", "gallery") + .describe("stage", "30 FPS latest-frame sampling deadline") + .describe("frame_rate_fps", "30"); + auto encode = media->add("FFmpeg.H264.encode", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.encode_media_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + encode.describe("owner", "gallery") + .describe("stage", "FFmpeg H.264 encode after Plot pixel publish") + .describe("backend", "FFmpeg") + .describe("codec", "H.264") + .describe("execution", "Taskflow worker"); + auto publish_ffmpeg = media->add("gallery.ffmpeg.webrtc.publish", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.publish_ffmpeg_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + publish_ffmpeg.describe("owner", "gallery") + .describe("transport", "WebRTC") + .describe("stage", "H.264 access unit enqueue"); + auto pack_pixels = media->add("gallery.websocket_pixels.pack", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.pack_pixel_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + pack_pixels.describe("owner", "gallery") + .describe("transport", "WebSocket") + .describe("pixel_conversion", "none") + .describe("stage", "native pixel protocol framing"); + auto publish_pixels = media->add( + "gallery.websocket_pixels.publish", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.publish_pixel_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + publish_pixels.describe("owner", "gallery") + .describe("transport", "WebSocket binary") + .describe("stage", "native pixel frame enqueue"); + auto complete = media->add("gallery.sample.complete", [weak] { + if (const auto owner = weak.lock()) { + auto& owner_data = static_cast(*owner->d); + try { owner_data.complete_media_frame(); } + catch (...) { owner_data.fail(std::current_exception()); } + } + }); + complete.describe("owner", "gallery") + .describe("stage", "sample metrics and ownership release"); + sample.precede(encode); + encode.precede(publish_ffmpeg); + 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)); + for (std::size_t slot = 0; slot < data.sources.size(); ++slot) { auto& source = data.sources[slot]; source.stream = source.entry.plot->subscribe( @@ -331,69 +480,17 @@ void Gallery_Video_Stream::bind_plots() { if (!owner) return; auto& owner_data = static_cast(*owner->d); if (owner_data.stopping.load(std::memory_order_acquire)) return; - /* - * 每个 Plot 只替换自己槽位的不可变最近帧快照。这里不能把 - * 每个源帧重新投递到编码任务,否则 24 个高帧率图 - * 会先制造任务洪泛,再让图集编码永远追赶旧任务。 - */ + /* 每个 Plot 只替换自己槽位的不可变最近帧快照。 */ owner_data.accept_frame(slot, std::move(frame)); }); - source.entry.plot->configure_stream(source.stream, tile_width, tile_height); + source.entry.plot->configure_stream( + source.stream, tile_width, tile_height); } - if (data.sources.empty()) throw std::invalid_argument("gallery video stream requires a Plot"); - data.encode_source_slot = data.sources.size() - 1; - auto media = std::make_unique("gallery.video.frame"); - auto compose = media->add("gallery.atlas.compose", [weak] { - if (const auto owner = weak.lock()) { - auto& owner_data = static_cast(*owner->d); - try { - owner_data.compose_media_frame(); - } - catch (...) { - owner_data.fail(std::current_exception()); - } - } - }); - compose.describe("owner", "gallery") - .describe("stage", "latest completed Plot frames to atlas"); - auto encode = media->add("FFmpeg.H264.encode", [weak] { - if (const auto owner = weak.lock()) { - auto& owner_data = static_cast(*owner->d); - try { - owner_data.encode_media_frame(); - } - catch (...) { - owner_data.fail(std::current_exception()); - } - } - }); - encode.describe("owner", "gallery") - .describe("stage", "FFmpeg H.264 encode in Scene completion pipeline") - .describe("backend", "FFmpeg") - .describe("codec", "H.264") - .describe("execution", "Taskflow worker"); - auto publish = media->add("gallery.webrtc.publish", [weak] { - if (const auto owner = weak.lock()) { - auto& owner_data = static_cast(*owner->d); - try { - owner_data.publish_media_frame(); - } - catch (...) { - owner_data.fail(std::current_exception()); - } - } - }); - publish.describe("owner", "gallery") - .describe("stage", "WebRTC frame delivery"); - compose.precede(encode); - encode.precede(publish); - data.sources[data.encode_source_slot].entry.plot->attach_scene_completion( - std::move(media)); } Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( - std::string connection, Stream_Handler handler, - Video_Readiness video_readiness, - Video_Diagnostics video_diagnostics) { + std::string connection, Gallery_Transport_Mode transport, + Stream_Handler handler, Transport_Readiness readiness, + Transport_Diagnostics diagnostics) { auto& data = static_cast(*d); if (data.stopping.load(std::memory_order_acquire)) throw std::logic_error("gallery video stream is shutting down"); if (data.failed.load(std::memory_order_acquire)) throw std::logic_error("gallery video stream is unavailable"); @@ -401,19 +498,19 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( if (connection.empty()) throw std::invalid_argument( "gallery video subscription requires a connection identity"); - if (!video_readiness) + if (!readiness) throw std::invalid_argument( - "gallery video subscription requires WebRTC readiness"); - if (!video_diagnostics) + "gallery media subscription requires transport readiness"); + if (!diagnostics) throw std::invalid_argument( - "gallery video subscription requires WebRTC diagnostics"); + "gallery media subscription requires transport diagnostics"); const auto id = data.next_consumer_id.fetch_add( 1, std::memory_order_relaxed); const Private::Consumer consumer{ - id, std::move(connection), + id, std::move(connection), transport, std::make_shared(std::move(handler)), - std::make_shared(std::move(video_readiness)), - std::make_shared(std::move(video_diagnostics))}; + std::make_shared(std::move(readiness)), + std::make_shared(std::move(diagnostics))}; auto current = data.consumers.load(std::memory_order_acquire); for (;;) { auto replacement = std::make_shared(); @@ -440,7 +537,7 @@ void Gallery_Video_Stream::request_video_key_frame() { auto& data = static_cast(*d); if (!data.stopping.load(std::memory_order_acquire) && !data.failed.load(std::memory_order_acquire)) - data.key_frame_requested.store(true, std::memory_order_release); + data.ffmpeg_transport.request_key_frame(); } void Gallery_Video_Stream::shutdown() noexcept { auto& data = static_cast(*d); @@ -456,15 +553,19 @@ void Gallery_Video_Stream::shutdown() noexcept { data.consumers.store(std::make_shared(), std::memory_order_release); } -std::string Gallery_Video_Stream::layout_description() const { +std::string Gallery_Video_Stream::layout_description( + Gallery_Transport_Mode transport) const { const auto& data = static_cast(*d); - const auto description = data.atlas->describe(); + const auto description = data.sampler->describe(); nlohmann::json plots = nlohmann::json::object(); for (const auto& source : description.sources) plots[source.id] = {{"column", source.column}, {"row", source.row}}; return nlohmann::json{ {"kind", "gallery_layout"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, + {"transport", transport == Gallery_Transport_Mode::ffmpeg + ? "ffmpeg" : "websocket_pixels"}, + {"frame_rate_fps", gallery_frame_rate}, {"columns", description.columns}, {"rows", description.rows}, {"tile_width", description.tile_width}, @@ -478,6 +579,12 @@ std::string Gallery_Video_Stream::layout_description() const { {"transport", "WebRTC"} } }, + {"pixels", { + {"transport", "WebSocket binary"}, + {"header_bytes", 48}, + {"payload", "native"}, + {"format_source", "binary_frame_header"} + }}, {"plots", std::move(plots)} }.dump(); } @@ -485,7 +592,7 @@ nlohmann::json Gallery_Video_Stream::diagnostics() const { nlohmann::json output{ {"kind", "gallery_metrics"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"sources", nlohmann::json::object()} }; static_cast*>(this) @@ -496,9 +603,10 @@ nlohmann::json Gallery_Video_Stream::diagnostics() const { const auto current = data.consumers.load(std::memory_order_acquire); output["consumer_count"] = current->size(); for (const auto& consumer : *current) { - if (!consumer.video_diagnostics) continue; - output["webrtc"] = (*consumer.video_diagnostics)(); - break; + if (!consumer.diagnostics) continue; + const auto key = consumer.transport == Gallery_Transport_Mode::ffmpeg + ? "webrtc" : "websocket_pixels"; + if (!output.contains(key)) output[key] = (*consumer.diagnostics)(); } return output; } diff --git a/web_server/src/Gallery_Video_Stream.hpp b/web_server/src/Gallery_Video_Stream.hpp index c47b487..137be5a 100644 --- a/web_server/src/Gallery_Video_Stream.hpp +++ b/web_server/src/Gallery_Video_Stream.hpp @@ -1,18 +1,24 @@ #pragma once -#include "H264_Encoder.hpp" #include "Plot.hpp" +#include "frame_sampling/ffmpeg/FFmpeg_Frame_Transport.hpp" +#include "frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.hpp" #include #include #include #include #include -#include #include #include namespace aethera::web { +enum struct Gallery_Transport_Mode : std::uint8_t { + ffmpeg, + websocket_pixels +}; + struct Gallery_Stream_Frame { - std::optional video{}; /* 页面唯一 H.264 图集 access unit;交付时转移所有权。 */ + std::shared_ptr video{}; /* FFmpeg 订阅者共享的 H.264 access unit。 */ + std::shared_ptr pixels{}; /* 像素订阅者共享的原生图集包。 */ std::string notification{}; /* 仅承载终止错误等低频控制通知。 */ }; @@ -31,8 +37,8 @@ struct Gallery_Video_Stream : Def, using Stream_Id = std::uint64_t; using Stream_Handler = std::function; - using Video_Readiness = std::function; - using Video_Diagnostics = std::function; + using Transport_Readiness = std::function; + using Transport_Diagnostics = std::function; Gallery_Video_Stream(); ~Gallery_Video_Stream(); @@ -41,13 +47,15 @@ struct Gallery_Video_Stream : Def, [[nodiscard]] static std::shared_ptr create( std::vector plots); [[nodiscard]] Stream_Id subscribe(std::string connection, + Gallery_Transport_Mode transport, Stream_Handler handler, - Video_Readiness video_readiness, - Video_Diagnostics video_diagnostics); + Transport_Readiness readiness, + Transport_Diagnostics diagnostics); void unsubscribe(Stream_Id stream); void request_video_key_frame(); void shutdown() noexcept; - [[nodiscard]] std::string layout_description() const; + [[nodiscard]] std::string layout_description( + Gallery_Transport_Mode transport) const; [[nodiscard]] nlohmann::json diagnostics() const; private: void bind_plots(); diff --git a/web_server/src/Gallery_Video_Stream.ipp b/web_server/src/Gallery_Video_Stream.ipp index 782f691..49178e1 100644 --- a/web_server/src/Gallery_Video_Stream.ipp +++ b/web_server/src/Gallery_Video_Stream.ipp @@ -1,5 +1,5 @@ #pragma once -#include "detail/Gallery_Frame_Atlas.hpp" +#include "frame_sampling/Frame_Sampler.hpp" #include #include #include @@ -13,37 +13,42 @@ struct Gallery_Video_Stream::Private : Prev_Private { struct Consumer { Stream_Id id{}; /* 当前网页媒体连接的订阅身份。 */ std::string connection{}; /* 同一浏览器页重连时保持稳定的协议身份。 */ + Gallery_Transport_Mode transport{Gallery_Transport_Mode::ffmpeg}; /* 当前连接选择的唯一媒体模型。 */ std::shared_ptr handler{}; /* 发布期间保持回调存活。 */ - std::shared_ptr video_readiness{}; /* 编码前查询发送窗口。 */ - std::shared_ptr video_diagnostics{}; /* WebRTC 独立诊断来源。 */ + std::shared_ptr readiness{}; /* 采样前查询对应传输窗口。 */ + std::shared_ptr diagnostics{}; /* 连接自己的传输诊断来源。 */ }; using Consumers = std::vector; Object* object{}; /* Def 机制所属最终对象。 */ - Sliding_Statistics compose_ms{600}; /* 图集合成统计计算器。 */ + Sliding_Statistics compose_ms{600}; /* 图集采样与变更槽复制统计计算器。 */ + Sliding_Statistics sample_delay_ms{600}; /* 固定 deadline 到实际采样的调度延迟。 */ Sliding_Statistics encode_ms{600}; /* H.264 编码统计计算器。 */ - Sliding_Statistics publish_ms{600}; /* WebRTC 发布统计计算器。 */ + Sliding_Statistics pixel_pack_ms{600}; /* 原生像素二进制封装统计计算器。 */ + Sliding_Statistics publish_ms{600}; /* 两种传输发布统计计算器。 */ std::vector sources{}; /* 已按业务标识排序的稳定图集来源。 */ - std::unique_ptr atlas{}; /* 最近完成帧与 RGBA 图集的唯一状态源。 */ - H264_Encoder encoder; /* 页面唯一 H.264 编码器。 */ - Plot_Render_Tick active_encode_tick{}; /* 当前媒体 DAG 的输入时钟。 */ - std::optional active_composition{}; /* 当前合成结果所有权。 */ - std::optional active_video{}; /* 当前编码结果所有权。 */ + std::unique_ptr sampler{}; /* latest 图像与 30 FPS deadline 的唯一状态源。 */ + frame_sampling::ffmpeg::FFmpeg_Frame_Transport ffmpeg_transport; /* FFmpeg 子模型的唯一编码上下文。 */ + std::optional active_sample{}; /* 当前 deadline 取得的原生图集视图。 */ + std::shared_ptr active_video{}; /* 当前编码结果不可变共享所有权。 */ + std::shared_ptr active_pixels{}; /* 当前原生像素协议包不可变共享所有权。 */ std::atomic_uint64_t next_consumer_id{1}; /* 订阅身份生成器。 */ std::atomic> consumers{ std::make_shared()}; /* 不可变订阅集合;写入使用 COW,发布只做一次原子读取。 */ std::atomic_bool stopping{}; /* 关闭开始后拒绝新工作。 */ - std::atomic_bool key_frame_requested{}; /* 下一帧关键帧请求。 */ - std::atomic_uint64_t skipped_encode_ticks{}; /* 被最新帧策略合并的编码 tick。 */ + std::atomic_uint64_t skipped_sample_ticks{}; /* 未到 deadline 或无可写消费者的采样触发数。 */ std::atomic_uint64_t delivered_clock_ticks{}; /* 已分发给 Plot 的时钟 tick。 */ std::atomic_bool failed{}; /* 首次 Unknown Failure 后停止热路径。 */ std::optional terminal_failure{}; /* 首次终止失败文本。 */ std::uint64_t encoded_frame_count{}; /* 累计编码成功帧数。 */ + std::uint64_t pixel_frame_count{}; /* 累计封装成功的原始像素帧数。 */ std::chrono::steady_clock::time_point metric_started{}; /* 当前指标窗口起点。 */ std::uint64_t metric_encoded_start{}; /* 窗口起点累计编码帧数。 */ + std::uint64_t metric_sample_start{}; /* 窗口起点累计采样帧数。 */ 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{}; /* 窗口起点累计像素帧数。 */ Private(); ~Private(); template @@ -52,19 +57,26 @@ struct Gallery_Video_Stream::Private : Prev_Private { object = static_cast(attached); } void initialize(std::vector plots); - [[nodiscard]] bool consumer_accepts_video() const noexcept; + [[nodiscard]] bool consumer_accepts( + Gallery_Transport_Mode transport) const noexcept; [[nodiscard]] bool remove_consumer(Stream_Id stream) noexcept; - void publish(std::optional video, + void publish(Gallery_Transport_Mode transport, + std::shared_ptr video, + std::shared_ptr pixels, std::string notification) noexcept; void fail(std::exception_ptr failure) noexcept; void accept_frame(std::size_t slot, std::shared_ptr frame); void update_metrics(const Plot_Render_Tick& tick, - const detail::Gallery_Atlas_Composition& composition, + const frame_sampling::Sampled_Gallery_Frame& sample, std::size_t encoded_bytes, - Video_Encoder_Backend encoder_backend); - void compose_media_frame(); + std::size_t pixel_bytes, + std::optional encoder_backend); + void sample_media_frame(); void encode_media_frame(); - void publish_media_frame(); + void pack_pixel_frame(); + void publish_ffmpeg_frame(); + void publish_pixel_frame(); + void complete_media_frame(); }; } diff --git a/web_server/src/Gallery_WebSocket.cpp b/web_server/src/Gallery_WebSocket.cpp index 1376523..7a3cc14 100644 --- a/web_server/src/Gallery_WebSocket.cpp +++ b/web_server/src/Gallery_WebSocket.cpp @@ -30,6 +30,7 @@ struct Gallery_WebSocket::Private { std::unique_ptr video; /* 当前浏览器连接的 WebRTC 发送会话。 */ Gallery_Video_Stream::Stream_Id subscription{}; /* 共享图集输出的订阅标识。 */ std::string connection_id{}; /* 浏览器页内重连稳定、跨页不同。 */ + Gallery_Transport_Mode transport{Gallery_Transport_Mode::ffmpeg}; /* 连接建立时协商的唯一媒体模型。 */ std::atomic_bool attached{}; /* start 成功后为真,并保证 close 只执行一次。 */ std::atomic_bool media_failed{}; /* 首次 WebRTC 发送 Unknown Failure 后停止重复报告。 */ }; @@ -37,11 +38,13 @@ struct Gallery_WebSocket::Private { Gallery_WebSocket::Gallery_WebSocket( drogon::WebSocketConnectionPtr connection, std::shared_ptr stream, - std::string connection_id) + std::string connection_id, + Gallery_Transport_Mode transport) : d(std::make_unique()) { d->connection = std::move(connection); d->stream = std::move(stream); d->connection_id = std::move(connection_id); + d->transport = transport; } Gallery_WebSocket::~Gallery_WebSocket() { close(); } @@ -49,33 +52,49 @@ Gallery_WebSocket::~Gallery_WebSocket() { close(); } void Gallery_WebSocket::start() { if (d->attached.exchange(true, std::memory_order_acq_rel)) return; const auto weak = weak_from_this(); - d->video = std::make_unique( - [weak](std::string signal) { - const auto socket = weak.lock(); - if (!socket) return; - const auto connection = socket->d->connection.lock(); - if (connection && connection->connected()) - connection->send(std::move(signal), drogon::WebSocketMessageType::Text); - }, - [weak] { - if (const auto socket = weak.lock()) - socket->d->stream->request_video_key_frame(); - }, - [weak](std::exception_ptr failure) { - if (const auto socket = weak.lock()) - socket->fail_media(std::move(failure)); - }); - d->subscription = d->stream->subscribe(d->connection_id, + if (d->transport == Gallery_Transport_Mode::ffmpeg) { + d->video = std::make_unique( + [weak](std::string signal) { + const auto socket = weak.lock(); + if (!socket) return; + const auto connection = socket->d->connection.lock(); + if (connection && connection->connected()) + connection->send(std::move(signal), + drogon::WebSocketMessageType::Text); + }, + [weak] { + if (const auto socket = weak.lock()) + socket->d->stream->request_video_key_frame(); + }, + [weak](std::exception_ptr failure) { + if (const auto socket = weak.lock()) + socket->fail_media(std::move(failure)); + }); + } + d->subscription = d->stream->subscribe(d->connection_id, d->transport, [weak](Gallery_Stream_Frame frame) { if (const auto socket = weak.lock()) socket->deliver(std::move(frame)); }, [weak] { const auto socket = weak.lock(); - return socket && socket->d->attached.load(std::memory_order_acquire) && - socket->d->video && socket->d->video->can_accept_video(); + if (!socket || + !socket->d->attached.load(std::memory_order_acquire)) + return false; + if (socket->d->transport == Gallery_Transport_Mode::websocket_pixels) { + const auto connection = socket->d->connection.lock(); + return connection && connection->connected(); + } + return socket->d->video && + socket->d->video->can_accept_video(); }, [weak] { const auto socket = weak.lock(); + if (socket && + socket->d->transport == Gallery_Transport_Mode::websocket_pixels) + return nlohmann::json{ + {"kind", "websocket_pixel_transport_state"}, + {"available", socket->d->attached.load( + std::memory_order_acquire)}}; return socket && socket->d->video ? socket->d->video->diagnostics() : nlohmann::json{{"kind", "webrtc_transport_state"}, @@ -83,9 +102,9 @@ void Gallery_WebSocket::start() { }); if (const auto connection = d->connection.lock(); connection && connection->connected()) - connection->send(d->stream->layout_description(), + connection->send(d->stream->layout_description(d->transport), drogon::WebSocketMessageType::Text); - d->video->start(); + if (d->video) d->video->start(); } void Gallery_WebSocket::deliver( @@ -95,11 +114,22 @@ void Gallery_WebSocket::deliver( if (!connection || !connection->connected()) return; if (frame.video && !d->media_failed.load(std::memory_order_acquire)) { try { - static_cast(d->video->send(std::move(*frame.video))); + static_cast(d->video->send(*frame.video)); } catch (...) { fail_media(std::current_exception()); } } + if (frame.pixels && + d->transport == Gallery_Transport_Mode::websocket_pixels && + !d->media_failed.load(std::memory_order_acquire)) { + try { + connection->send(frame.pixels->bytes, + drogon::WebSocketMessageType::Binary); + } + catch (...) { + fail_media(std::current_exception()); + } + } if (!frame.notification.empty()) { try { connection->send(std::move(frame.notification), @@ -117,7 +147,7 @@ void Gallery_WebSocket::fail_media(std::exception_ptr failure) noexcept { connection->send(nlohmann::json{ {"kind", "gallery_error"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"message", media_failure_description(failure)}}.dump(), drogon::WebSocketMessageType::Text); } @@ -129,15 +159,15 @@ void Gallery_WebSocket::receive(std::string_view message) { if (json.is_discarded() || !json.is_object()) return; const auto kind = json.value("kind", std::string{}); try { - if (kind == "webrtc_answer") { + if (kind == "webrtc_answer" && d->video) { const auto sdp = json.value("sdp", std::string{}); if (!sdp.empty()) d->video->accept_answer(sdp); - } else if (kind == "webrtc_candidate") { + } else if (kind == "webrtc_candidate" && d->video) { const auto candidate = json.value("candidate", std::string{}); if (!candidate.empty()) d->video->add_remote_candidate( candidate, json.value("mid", std::string{})); - } else if (kind == "request_key_frame") { + } else if (kind == "request_key_frame" && d->video) { d->stream->request_video_key_frame(); } } @@ -178,7 +208,7 @@ void Gallery_WebSocket_Controller::handleNewConnection( connection->send(nlohmann::json{ {"kind", "gallery_error"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"message", "unknown gallery media group"}}.dump(), drogon::WebSocketMessageType::Text); } @@ -195,8 +225,15 @@ void Gallery_WebSocket_Controller::handleNewConnection( if (connection_id.empty()) throw std::invalid_argument( "gallery WebSocket requires a browser connection identity"); + const auto transport_name = request->getParameter("transport"); + Gallery_Transport_Mode transport{Gallery_Transport_Mode::ffmpeg}; + if (transport_name == "websocket_pixels") + transport = Gallery_Transport_Mode::websocket_pixels; + else if (!transport_name.empty() && transport_name != "ffmpeg") + throw std::invalid_argument( + "unknown gallery transport; expected ffmpeg or websocket_pixels"); auto socket = std::make_shared( - connection, stream, connection_id); + connection, stream, connection_id, transport); connection->setContext(socket); connection->setPingMessage( "aethera-gallery-video", std::chrono::seconds(20)); @@ -209,7 +246,7 @@ void Gallery_WebSocket_Controller::handleNewConnection( connection->send(nlohmann::json{ {"kind", "gallery_error"}, {"protocol", "aethera.gallery.video"}, - {"version", 2}, + {"version", 3}, {"message", media_failure_description( std::current_exception())}}.dump(), drogon::WebSocketMessageType::Text); diff --git a/web_server/src/Gallery_WebSocket.hpp b/web_server/src/Gallery_WebSocket.hpp index 9d8fe13..9b37d38 100644 --- a/web_server/src/Gallery_WebSocket.hpp +++ b/web_server/src/Gallery_WebSocket.hpp @@ -11,7 +11,8 @@ struct Gallery_WebSocket final public: Gallery_WebSocket(drogon::WebSocketConnectionPtr connection, std::shared_ptr stream, - std::string connection_id); + std::string connection_id, + Gallery_Transport_Mode transport); ~Gallery_WebSocket(); Gallery_WebSocket(const Gallery_WebSocket&) = delete; Gallery_WebSocket& operator=(const Gallery_WebSocket&) = delete; diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index 874fc4f..a8968dc 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -8,6 +8,7 @@ #include #include #include +#include #include #include #include @@ -18,6 +19,7 @@ #include #include #include +#include #include #include #include @@ -101,7 +103,7 @@ nlohmann::json Frame_Policy::schema() const { {"technical_description", "Authoritative per-scene periodic render switch."}, {"value", current.render_enabled}}, {{"key", "video_enabled"}, {"label", "图集视频传输"}, {"editor", "boolean"}, - {"editable", true}, {"description", "控制完成帧是否进入页面级 RGBA 图集;默认开启,用于完整链路压测。"}, + {"editable", true}, {"description", "控制完成帧是否进入页面级采样器;2D BGRA 与 3D RGBA 均保持原生格式。"}, {"technical_description", "Authoritative tile publication switch for the shared gallery video."}, {"value", current.video_enabled}}, {{"key", "pacing_mode"}, {"label", "服务端帧策略"}, {"editor", "select"}, @@ -319,8 +321,8 @@ struct Plot::Private { enum struct Frame_State : std::uint8_t { available, - in_flight, - callback_retired + rendering, /* Scene::advance -> Plot pixel publish,不可重入。 */ + consuming /* 外接 Taskflow 正在消费已发布帧;允许下一帧渲染。 */ }; struct Managed_Frame { @@ -353,16 +355,17 @@ struct Plot::Private { Frame_Policy frame_policy{}; Frame_Scheduler::Timer frame_timer{}; /* 每 Plot/Scene 只有轻量时间轮节点,不持有线程。 */ static constexpr std::size_t scene_frame_capacity{3}; - std::array frame_slots{}; /* Scene 借用的稳定三缓冲物理帧。 */ - std::atomic_size_t callback_retired_slot{scene_frame_capacity}; + std::array frame_slots{}; /* Scene 与外接消费者共享生命周期的稳定三缓冲。 */ Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */ std::atomic> pending_tick{}; std::atomic_bool tick_task_scheduled{}; /* 唯一短任务准入;不占用 Worker 等待。 */ + std::atomic_bool render_admission_busy{}; /* view->update/Scene::advance 到 pixel publish 的唯一准入门。 */ std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()}; std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */ std::atomic_uint64_t coalesced_tick_count{}; /* 尚未消费时被更新 tick 替换的旧 tick 总数。 */ std::atomic_uint64_t policy_skip_count{}; /* 帧策略拒绝的 tick 总数。 */ - std::atomic_uint64_t preparation_busy_count{}; /* 本 Plot Prepare 已在执行而跳过的提交次数。 */ + std::atomic_uint64_t 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 的帧总数。 */ @@ -373,8 +376,19 @@ struct Plot::Private { std::atomic_uint64_t taskflow_trace_control{}; std::array>, maximum_taskflow_trace_frames> taskflow_trace_slots{}; - Task_Node completion_tail{}; /* Scene 完成图中当前最后一个业务阶段。 */ - std::vector> completion_extensions{}; /* 生命周期覆盖 Scene 对子图的借用。 */ + 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_Node completion_tail{}; /* Scene 图内固定停在 plot.frame.publish。 */ + 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}; + Frame_Statistics_State completed_frame_statistics_state{}; + mutable std::mutex completed_frame_statistics_mutex{}; template Private(std::unique_ptr value_scene, @@ -400,14 +414,28 @@ 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 arm_tick_consumer(std::weak_ptr lifetime); + void release_render_admission(std::weak_ptr lifetime); void consume_tick(std::weak_ptr lifetime); void refresh_schedule(); void clock_tick(const Plot_Render_Tick& tick); void render_frame(Plot_Render_Tick tick); void publish_completed_frame(); + void consume_completed_frame(Render_Frame* frame); void retire_completed_frame(Render_Frame* 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, + const Taskflow_Frame_Trace& trace); + [[nodiscard]] nlohmann::json trace_response( + const std::atomic_uint64_t& control, + const std::atomic_size_t& remaining, + const std::array>, + maximum_taskflow_trace_frames>& slots) const; void fail(std::exception_ptr failure) noexcept; }; @@ -492,25 +520,64 @@ void Plot::Private::refresh_schedule() { else frame_timer.start_periodic(fps); } -void Plot::Private::consume_tick(std::weak_ptr lifetime) { - const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel); - if (tick) clock_tick(*tick); - tick_task_scheduled.store(false, std::memory_order_release); - if (pending_tick.load(std::memory_order_acquire) && - !tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) { - aethera::schedule_task("web.plot.tick.consume", [lifetime] { - const auto plot = lifetime.lock(); - if (!plot) return; - try { plot->d->consume_tick(lifetime); } - catch (...) { plot->d->fail(std::current_exception()); } - }); +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; + } + 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; + } } } +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 (tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return; + aethera::schedule_task("web.plot.tick.consume", [lifetime] { + const auto plot = lifetime.lock(); + if (!plot) return; + try { plot->d->consume_tick(lifetime); } + catch (...) { plot->d->fail(std::current_exception()); } + }); +} + +void Plot::Private::release_render_admission(std::weak_ptr lifetime) { + if (!render_admission_busy.exchange(false, std::memory_order_acq_rel)) + return; + if (pending_tick.load(std::memory_order_acquire)) + deferred_resume_count.fetch_add(1, std::memory_order_relaxed); + 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); + if (tick) clock_tick(*tick); + } + tick_task_scheduled.store(false, std::memory_order_release); + 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 (!frame_policy.accept_periodic_tick(tick.time_milliseconds)) { + if (tick.sequence != 0 && + !frame_policy.accept_periodic_tick(tick.time_milliseconds)) { policy_skip_count.fetch_add(1, std::memory_order_relaxed); return; } @@ -530,27 +597,85 @@ 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>, + maximum_taskflow_trace_frames>& slots, + const Taskflow_Frame_Trace& value) { + auto trace = std::make_shared(value); + auto state = control.load(std::memory_order_acquire); + for (;;) { + const auto requested = static_cast(state >> 32U); + const auto captured = static_cast(state); + if (captured >= requested) return; + slots[captured].store(trace, std::memory_order_release); + const auto next = (static_cast(requested) << 32U) | + static_cast(captured + 1U); + if (control.compare_exchange_weak( + state, next, std::memory_order_release, + std::memory_order_acquire)) + return; + } +} + +nlohmann::json Plot::Private::trace_response( + const std::atomic_uint64_t& control, + const std::atomic_size_t& remaining, + const std::array>, + maximum_taskflow_trace_frames>& slots) const { + nlohmann::json frames = nlohmann::json::array(); + const auto state = 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 = slots[index].load(std::memory_order_acquire)) + frames.push_back(taskflow_trace_json(*trace)); + const auto left = remaining.load(std::memory_order_acquire); + return { + {"protocol", "aethera.taskflow.frames"}, {"version", 1}, + {"requested", requested}, {"remaining", left}, + {"captured", frames.size()}, + {"complete", requested != 0 && frames.size() == requested}, + {"frames", std::move(frames)}}; +} + void Plot::Private::render_frame(Plot_Render_Tick tick) { if (terminal_failure.load(std::memory_order_acquire)) return; const auto streams = stream_snapshot(); const auto pacing = frame_policy.read(); if (!pacing.render_enabled || streams.consumers->empty()) return; - std::size_t slot_index{}; - Managed_Frame* managed{}; - /* Frame_State 是 Plot 数据准备生命周期的唯一权威来源。槽位通过 CAS - * 准入,完成图和时钟任务不再共享一把外围锁。 */ - if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) { - return slot.state.load(std::memory_order_acquire) == - Frame_State::in_flight; - })) { + bool admission_expected = false; + if (!render_admission_busy.compare_exchange_strong( + admission_expected, true, std::memory_order_acq_rel, + std::memory_order_acquire)) { preparation_busy_count.fetch_add(1, std::memory_order_relaxed); + defer_tick(tick); return; } + + std::size_t slot_index{}; + Managed_Frame* managed{}; + /* + * 只有 rendering 槽受 Scene 不可重入门约束;consuming 槽表示上一帧 + * 已经完成 Plot 像素发布,外接 H264/WebRTC 仍可继续持有该物理帧的 + * 诊断生命周期。只要还有 available 槽,下一帧即可进入。 + */ for (std::size_t index = 0; index < frame_slots.size(); ++index) { auto expected = Frame_State::available; if (!frame_slots[index].state.compare_exchange_strong( - expected, Frame_State::in_flight, + expected, Frame_State::rendering, std::memory_order_acq_rel, std::memory_order_acquire)) continue; slot_index = index; @@ -559,6 +684,14 @@ 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); + /* + * 三个槽都仍被外接消费者持有时,只保留 latest pending。这里绝不能 + * 立即 arm tick consumer,否则会在没有任何槽可用期间形成 + * consume -> no slot -> consume 的 Taskflow 任务风暴。真正的唤醒点 + * 是 retire_completed_frame:某个 consuming 槽变回 available 后只唤醒一次。 + */ + render_admission_busy.store(false, std::memory_order_release); return; } managed->presentation_time = @@ -566,7 +699,7 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { std::chrono::duration(tick.time_milliseconds)); const auto rollback_unsubmitted = [this, slot_index] { auto& slot = frame_slots[slot_index]; - auto expected = Frame_State::in_flight; + auto expected = Frame_State::rendering; static_cast(slot.state.compare_exchange_strong( expected, Frame_State::available, std::memory_order_acq_rel, std::memory_order_acquire)); @@ -616,6 +749,7 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { scene_rejection_count.fetch_add(1, std::memory_order_relaxed); rollback_unsubmitted(); restore_taskflow_trace_claim(); + release_render_admission(lifetime); } else { taskflow_trace_claimed = false; frame_policy.frame_submitted(); @@ -643,12 +777,14 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { scene_rejection_count.fetch_add(1, std::memory_order_relaxed); rollback_unsubmitted(); restore_taskflow_trace_claim(); + release_render_admission(lifetime); if (result == Render_Scene_3D::Render_Result::backend_unavailable) throw std::runtime_error("3D render backend became unavailable before submission"); } catch (...) { rollback_unsubmitted(); restore_taskflow_trace_claim(); + release_render_admission(lifetime); throw; } } @@ -657,29 +793,19 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { void Plot::Private::publish_completed_frame() { Render_Frame* frame{}; Managed_Frame* managed{}; - /* Scene 在用户回调返回之后才标记 callback_finished。下一次串行完成图 - * 开始时,上一 callback_retired 槽才可归还。 */ - const auto retired_index = callback_retired_slot.exchange( - scene_frame_capacity, std::memory_order_acq_rel); - if (retired_index != scene_frame_capacity) { - auto expected = Frame_State::callback_retired; - if (!frame_slots[retired_index].state.compare_exchange_strong( - expected, Frame_State::available, std::memory_order_acq_rel, - std::memory_order_acquire)) - throw std::logic_error("Plot retired frame state is inconsistent"); - } for (std::size_t index = 0; index < frame_slots.size(); ++index) { if (frame_slots[index].state.load(std::memory_order_acquire) != - Frame_State::in_flight) continue; + Frame_State::rendering) + continue; if (managed) - throw std::logic_error("Plot has multiple active Scene frames"); + throw std::logic_error("Plot has multiple frames in Scene rendering"); frame = std::visit( [](const auto& value) -> Render_Frame* { return value.get(); }, frame_slots[index].frame); managed = &frame_slots[index]; } if (!managed) - throw std::logic_error("Scene completion graph has no active Plot frame"); + throw std::logic_error("Scene completion graph has no rendering Plot frame"); try { const auto pacing = frame_policy.read(); @@ -729,89 +855,144 @@ void Plot::Private::publish_completed_frame() { static_cast(std::max(0, std::chrono::duration_cast( std::chrono::steady_clock::now() - publish_started).count()))); + auto expected = Frame_State::rendering; + if (!managed->state.compare_exchange_strong( + expected, Frame_State::consuming, std::memory_order_acq_rel, + std::memory_order_acquire)) + throw std::logic_error("Plot frame left rendering before pixel publish"); } catch (...) { throw; } } -void Plot::Private::retire_completed_frame(Render_Frame* frame) { +void Plot::Private::consume_completed_frame(Render_Frame* frame) { if (!frame) throw std::invalid_argument("Plot received a null completed frame"); - std::size_t slot_index{scene_frame_capacity}; - for (std::size_t index = 0; index < frame_slots.size(); ++index) { + Managed_Frame* managed{}; + for (auto& slot : frame_slots) { auto* address = std::visit( [](const auto& value) -> Render_Frame* { return value.get(); }, - frame_slots[index].frame); + slot.frame); if (address != frame) continue; - auto expected = Frame_State::in_flight; - if (!frame_slots[index].state.compare_exchange_strong( - expected, Frame_State::callback_retired, - std::memory_order_acq_rel, std::memory_order_acquire)) - throw std::logic_error("completed Plot frame is not in flight"); - slot_index = index; + if (slot.state.load(std::memory_order_acquire) != Frame_State::consuming) + throw std::logic_error("completed Plot frame was not published"); + managed = &slot; break; } - if (slot_index == scene_frame_capacity) - throw std::logic_error( - "frame callback has no externally owned active frame"); - auto no_retired = scene_frame_capacity; - if (!callback_retired_slot.compare_exchange_strong( - no_retired, slot_index, std::memory_order_acq_rel, + 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 内执行,但不会阻塞下一帧渲染。 + */ + 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)) - throw std::logic_error("Plot has more than one callback-retired frame"); - if (frame->taskflow_trace_requested()) { - auto trace = std::make_shared( - frame->taskflow_trace()); - auto control = taskflow_trace_control.load(std::memory_order_acquire); - for (;;) { - const auto requested = static_cast(control >> 32U); - const auto captured = static_cast(control); - if (captured >= requested) break; - 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( - control, next, std::memory_order_release, - std::memory_order_acquire)) - break; + 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->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)); } - if (frame_policy.frame_completed()) { - const auto weak = lifetime; - aethera::schedule_task("web.plot.render.deferred", [weak] { - const auto owner = weak.lock(); - if (!owner) return; - try { - const auto now = std::chrono::steady_clock::now(); - const auto elapsed = now - owner->d->clock_origin; - owner->d->render_frame(Plot_Render_Tick{ - now, 0, - std::chrono::duration(elapsed).count()}); - } - catch (...) { - owner->d->fail(std::current_exception()); - } - }); + 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; } } +void Plot::Private::retire_completed_frame(Render_Frame* frame) { + if (!frame) + throw std::invalid_argument("Plot received a null retired frame"); + Managed_Frame* managed{}; + for (auto& slot : frame_slots) { + auto* address = std::visit( + [](const auto& value) -> Render_Frame* { return value.get(); }, + slot.frame); + if (address != frame) continue; + managed = &slot; + break; + } + if (!managed) + throw std::logic_error("retired frame has no owned Plot slot"); + + { + std::lock_guard lock(completed_frame_statistics_mutex); + completed_frame_statistics_state = completed_frame_statistics.submit( + *frame, std::holds_alternative>(managed->frame) + ? Frame_Dimension::two_dimensional + : Frame_Dimension::three_dimensional); + } + + if (frame->taskflow_trace_requested()) + store_trace(taskflow_trace_control, taskflow_trace_slots, + frame->taskflow_trace()); + + auto expected = Frame_State::consuming; + if (!managed->state.compare_exchange_strong( + expected, Frame_State::available, std::memory_order_acq_rel, + std::memory_order_acquire)) + throw std::logic_error("retired Plot frame is not consuming"); + + /* 三个消费者槽曾全部占满时,退役一个槽后继续 latest pending。 */ + arm_tick_consumer(lifetime); +} + void Plot::Private::attach_completion( std::unique_ptr completion) { if (!completion || completion->empty()) throw std::invalid_argument("Plot completion pipeline is empty"); - auto& scene_completion = std::visit( - [](auto& scene_value) -> Task_Graph& { - return scene_value->completion_taskflow(); - }, scene); completion_extensions.push_back(std::move(completion)); - auto extension = scene_completion.compose( + auto extension = post_publish_graph.compose( completion_extensions.back()->name(), *completion_extensions.back()); - extension.describe("owner", "scene") - .describe("stage", "frame pipeline extension"); - completion_tail.precede(extension); - completion_tail = std::move(extension); + 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, @@ -838,29 +1019,38 @@ void Plot::ensure_started() { tick.issued_at, tick.sequence, tick.time_milliseconds}); } }); - d->refresh_schedule(); /* - * Kernel Frame_Scheduler 只产生每个 Scene 独立的周期 tick; - * Scene::render(Frame*) 与外部 Frame 所有权保持原样。完成回调先归还 - * 当前帧;若期间出现 immediate 请求,则回调返回后重新投递 Taskflow。 + * 2D 的 callback 在 Scene render admission 已释放后运行:先执行所有 + * post-publish 外接 DAG;Scene 完成 frame_ready 与 trace 收口后,再由 + * retired callback 归还物理槽。这样 H264(N) 可与 Render(N+1) 重叠。 */ if (auto* scene = std::get_if>(&d->scene)) { - (*scene)->set_frame_callback( - [weak](Frame_2D* frame) { - if (auto owner = weak.lock()) { - try { owner->d->retire_completed_frame(frame); } - catch (...) { owner->d->fail(std::current_exception()); } - } - }); + (*scene)->set_frame_callback([weak](Frame_2D* frame) { + if (auto owner = weak.lock()) { + try { owner->d->consume_completed_frame(frame); } + catch (...) { owner->d->fail(std::current_exception()); } + } + }); + (*scene)->set_frame_retired_callback([weak](Frame_2D* frame) { + if (auto owner = weak.lock()) { + try { owner->d->retire_completed_frame(frame); } + catch (...) { owner->d->fail(std::current_exception()); } + } + }); } else { std::get>(d->scene)->set_frame_callback( [weak](Frame_3D* frame) { - if (auto owner = weak.lock()) { - try { owner->d->retire_completed_frame(frame); } - catch (...) { owner->d->fail(std::current_exception()); } + if (auto owner = weak.lock()) { + try { + owner->d->consume_completed_frame(frame); + owner->d->retire_completed_frame(frame); } + catch (...) { owner->d->fail(std::current_exception()); } + } }); } + /* callback 必须先于周期时钟安装,避免首帧在初始化窗口进入 Scene。 */ + d->refresh_schedule(); }); } @@ -935,41 +1125,17 @@ 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); - const auto next = std::make_shared(std::move(tick)); - if (d->pending_tick.exchange(next, std::memory_order_acq_rel)) - d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed); - if (d->tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return; - const auto weak = weak_from_this(); - aethera::schedule_task("web.plot.tick.consume", [weak] { - const auto owner = weak.lock(); - if (!owner) return; - try { - owner->d->consume_tick(weak); - } - catch (...) { - owner->d->fail(std::current_exception()); - } - }); + d->defer_tick(tick); + d->arm_tick_consumer(weak_from_this()); } void Plot::render_once() { ensure_started(); if (!d->frame_policy.request_immediate()) return; - const auto weak = weak_from_this(); - aethera::schedule_task("web.plot.render.once", [weak] { - const auto owner = weak.lock(); - if (!owner) return; - try { - const auto elapsed = std::chrono::steady_clock::now() - - owner->d->clock_origin; - owner->d->render_frame(Plot_Render_Tick{ - std::chrono::steady_clock::now(), - 0, std::chrono::duration(elapsed).count()}); - } - catch (...) { - owner->d->fail(std::current_exception()); - } - }); + 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()}); } void Plot::submit_input(Plot_Input_Event event) { @@ -1021,10 +1187,25 @@ nlohmann::json Plot::diagnostics() const { std::uint64_t dropped_sequences{}; double frame_rate{}; bool is_3d{}; - const auto read_statistics = [&](const auto& state) { - const auto& statistics = state.frame_statistics; - append_statistic_json(frame_statistics, statistics); + const auto read_scene_statistics = [&](const auto& state) { append_event_statistics_json(input_statistics, state.event_statistics); + }; + std::visit([&](const auto& scene) { + using Scene_Pointer = std::remove_cvref_t; + if constexpr (std::same_as>) { + scene->template access_state( + read_scene_statistics); + } else { + is_3d = true; + scene->template access_state( + read_scene_statistics); + } + }, d->scene); + + { + std::lock_guard lock(d->completed_frame_statistics_mutex); + const auto statistics = d->completed_frame_statistics_state; + append_statistic_json(frame_statistics, statistics); identity = statistics.identity; created_time_unix_ns = statistics.created_time_unix_ns; dropped_sequences = statistics.dropped_sequences; @@ -1032,18 +1213,7 @@ nlohmann::json Plot::diagnostics() const { static_cast(Frame_Statistic::frame_interval_ms)]; frame_rate = interval.trimmed_average > 0.0 ? 1'000.0 / interval.trimmed_average : 0.0; - }; - std::visit([&](const auto& scene) { - using Scene_Pointer = std::remove_cvref_t; - if constexpr (std::same_as>) { - scene->template access_state( - read_statistics); - } else { - is_3d = true; - scene->template access_state( - read_statistics); - } - }, d->scene); + } const auto pacing = d->frame_policy.read(); const auto stream = d->stream_snapshot(); @@ -1057,8 +1227,7 @@ nlohmann::json Plot::diagnostics() const { } const auto format = is_3d ? pixel_format_name(Frame_3D::native_pixel_format) - : pixel_format_name(pacing.video_enabled - ? render_2d::Pixel_Format::rgba8 : Frame_2D::native_pixel_format); + : pixel_format_name(Frame_2D::native_pixel_format); const auto native_format = is_3d ? pixel_format_name(Frame_3D::native_pixel_format) : pixel_format_name(Frame_2D::native_pixel_format); @@ -1090,6 +1259,8 @@ nlohmann::json Plot::diagnostics() const { {"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)}, {"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)}}}, @@ -1149,26 +1320,45 @@ void Plot::request_taskflow_trace(std::size_t frame_count) { } nlohmann::json Plot::taskflow_trace() const { - nlohmann::json frames = nlohmann::json::array(); - const auto control = d->taskflow_trace_control.load( - std::memory_order_acquire); - const auto requested = static_cast(control >> 32U); - const auto captured = static_cast(control); - for (std::uint32_t index = 0; index < captured; ++index) - if (const auto trace = d->taskflow_trace_slots[index].load( + return d->trace_response(d->taskflow_trace_control, + 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)) - frames.push_back(taskflow_trace_json(*trace)); - const auto remaining = d->taskflow_trace_remaining.load( - std::memory_order_acquire); - return { - {"protocol", "aethera.taskflow.frames"}, {"version", 1}, - {"requested", requested}, {"remaining", remaining}, - {"captured", frames.size()}, - {"complete", requested != 0 && frames.size() == requested}, - {"frames", std::move(frames)}}; + 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->reset_frame_statistics(); }, d->scene); + std::lock_guard lock(d->completed_frame_statistics_mutex); + d->completed_frame_statistics.reset(); + d->completed_frame_statistics_state = {}; } } diff --git a/web_server/src/Plot.hpp b/web_server/src/Plot.hpp index 8454123..d99f77b 100644 --- a/web_server/src/Plot.hpp +++ b/web_server/src/Plot.hpp @@ -91,6 +91,9 @@ 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; diff --git a/web_server/src/Web_Server.cpp b/web_server/src/Web_Server.cpp index 1623c70..b723956 100644 --- a/web_server/src/Web_Server.cpp +++ b/web_server/src/Web_Server.cpp @@ -101,6 +101,83 @@ nlohmann::json taskflow_runtime_json() { {"task_types", std::move(task_types)}, {"workers", std::move(workers)}}; } + +bool taskflow_graph_contains_gallery_media(const nlohmann::json& graph) { + if (!graph.is_object() || !graph.contains("nodes") || + !graph["nodes"].is_array()) + return false; + for (const auto& node : graph["nodes"]) { + if (!node.is_object()) continue; + const auto name = node.value("name", std::string{}); + if (name == "gallery.sample.capture" || + name == "FFmpeg.H264.encode" || + name == "gallery.ffmpeg.webrtc.publish" || + name == "gallery.websocket_pixels.pack" || + name == "gallery.websocket_pixels.publish" || + name == "gallery.sample.complete") + return true; + } + return false; +} + +/* + * Gallery 的媒体 DAG 实际只在组内唯一 encode-source Plot 上执行。诊断接口 + * 在服务端把该真实执行帧中由 Task_Graph 包装器捕获到的媒体 graph 附加到 + * 当前 Plot 的 trace 响应;不复制执行、不伪造 FFmpeg 节点,也不让前端知道 + * encode-source 是哪个 Plot。 + */ +nlohmann::json merge_gallery_media_trace(nlohmann::json plot_trace, + const nlohmann::json& media_trace) { + if (!plot_trace.is_object() || !media_trace.is_object() || + !plot_trace.contains("frames") || !plot_trace["frames"].is_array() || + !media_trace.contains("frames") || !media_trace["frames"].is_array()) + return plot_trace; + + auto& plot_frames = plot_trace["frames"]; + const auto& media_frames = media_trace["frames"]; + const auto count = std::min(plot_frames.size(), media_frames.size()); + + /* + * 只把真实媒体 graph 附加到已有 Plot trace;绝不裁剪 Plot 自己的帧, + * 也不改写 Plot trace 的 requested/captured/remaining/complete 语义。 + */ + for (std::size_t index = 0; index < count; ++index) { + auto& output = plot_frames[index]; + const auto& media = media_frames[index]; + if (!output.contains("graphs") || !output["graphs"].is_array() || + !output.contains("executions") || !output["executions"].is_array() || + !media.contains("graphs") || !media["graphs"].is_array() || + !media.contains("executions") || !media["executions"].is_array()) + continue; + + std::vector media_native_ids; + for (const auto& graph : media["graphs"]) { + if (!taskflow_graph_contains_gallery_media(graph)) continue; + auto appended = graph; + appended["stage"] = "gallery.media"; + for (const auto& node : graph["nodes"]) + if (node.is_object() && node.contains("native_id") && + node["native_id"].is_string()) + media_native_ids.push_back(node["native_id"].get()); + output["graphs"].push_back(std::move(appended)); + } + if (media_native_ids.empty()) continue; + for (const auto& execution : media["executions"]) { + if (!execution.is_object() || !execution.contains("native_id") || + !execution["native_id"].is_string()) + continue; + const auto native_id = execution["native_id"].get(); + if (std::ranges::find(media_native_ids, native_id) != + media_native_ids.end()) + output["executions"].push_back(execution); + } + output["gallery_media_sequence"] = media.value("sequence", 0ULL); + output["gallery_media_correlation_id"] = + media.value("correlation_id", 0ULL); + } + + return plot_trace; +} } int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) { @@ -146,12 +223,18 @@ 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>(); - const auto add_media_group = [&gallery_streams, &plot_media]( + auto plot_media_source = std::make_shared< + std::unordered_map>(); + const auto add_media_group = [&gallery_streams, &plot_media, &plot_media_source]( std::string id, std::vector entries) { + if (entries.empty()) return; const auto path = "/ws/gallery?group=" + id; - for (const auto& entry : entries) + 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); + } gallery_streams->emplace( std::move(id), Gallery_Video_Stream::create(std::move(entries))); }; @@ -237,7 +320,7 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) callback(json_response(taskflow_runtime_json())); }, {drogon::Get}); - app.registerHandler("/plot/{1}/taskflow", [plots]( + app.registerHandler("/plot/{1}/taskflow", [plots, plot_media_source]( const drogon::HttpRequestPtr& request, std::function&& callback, std::string plot_id) { @@ -246,6 +329,10 @@ 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 : find_plot(*plots, source_found->second); + if (!media_plot) media_plot = plot; try { if (request->method() == drogon::Post) { const auto input = nlohmann::json::parse(request->body()); @@ -255,9 +342,37 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) "Taskflow trace requires unsigned frame_count")); return; } - plot->request_taskflow_trace(input["frame_count"].get()); + const auto frame_count = input["frame_count"].get(); + const auto trace_active = [](const nlohmann::json& state) { + return state.value("requested", std::size_t{}) > + state.value("captured", std::size_t{}); + }; + if (trace_active(plot->taskflow_trace()) || + trace_active(media_plot->post_publish_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); + plot->request_taskflow_trace(frame_count); } - callback(json_response(plot->taskflow_trace())); + auto output = plot->taskflow_trace(); + const auto media = media_plot->post_publish_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); + output["media_requested"] = media.value("requested", 0U); + output["media_captured"] = media.value("captured", 0U); + output["media_remaining"] = media.value("remaining", 0U); + /* Scene 与异步媒体图分别完成捕获后才结束轮询。 */ + output["complete"] = plot_complete && media_complete; + output["remaining"] = std::max( + output.value("remaining", 0U), + media.value("remaining", 0U)); + callback(json_response(std::move(output))); } catch (const nlohmann::json::exception&) { callback(error_response(drogon::k400BadRequest, diff --git a/web_server/src/detail/Gallery_Frame_Atlas.cpp b/web_server/src/detail/Gallery_Frame_Atlas.cpp index 88ba2e3..6165d3d 100644 --- a/web_server/src/detail/Gallery_Frame_Atlas.cpp +++ b/web_server/src/detail/Gallery_Frame_Atlas.cpp @@ -171,6 +171,7 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { result.sources.push_back({ latest_completion->sequence, latest_completion->correlation_id, + latest_completion->presentation_time, latest_completion->rendered_sequence, latest_completion->rendered_correlation_id, published->completion_count, diff --git a/web_server/src/detail/Gallery_Frame_Atlas.hpp b/web_server/src/detail/Gallery_Frame_Atlas.hpp index 546f850..91dbfc8 100644 --- a/web_server/src/detail/Gallery_Frame_Atlas.hpp +++ b/web_server/src/detail/Gallery_Frame_Atlas.hpp @@ -1,5 +1,6 @@ #pragma once #include "../Plot.hpp" +#include #include #include #include @@ -25,6 +26,7 @@ struct Gallery_Atlas_Description { struct Gallery_Atlas_Source_Progress { std::uint64_t completion_sequence{}; /* 当前槽位最近一次逻辑完成回调的局部序号。 */ std::uint64_t completion_correlation_id{}; /* 最近逻辑回调对应的页面时钟序号。 */ + std::chrono::microseconds completion_presentation_time{}; /* 最近逻辑回调的页面媒体时间。 */ std::uint64_t rendered_sequence{}; /* 当前像素对应的真实 Scene/GPU 画面序号。 */ std::uint64_t rendered_correlation_id{}; /* 当前真实画面对应的页面时钟序号。 */ std::uint64_t completion_count{}; /* 本图集生命周期内接受的逻辑完成回调数。 */ diff --git a/web_server/src/frame_sampling/Frame_Sampler.cpp b/web_server/src/frame_sampling/Frame_Sampler.cpp new file mode 100644 index 0000000..e81e473 --- /dev/null +++ b/web_server/src/frame_sampling/Frame_Sampler.cpp @@ -0,0 +1,135 @@ +#include "Frame_Sampler.hpp" +#include +#include +#include +#include +#include + +namespace aethera::web::frame_sampling { +namespace { +using Clock = std::chrono::steady_clock; + +std::int64_t clock_nanoseconds(Clock::time_point value) noexcept { + return std::chrono::duration_cast( + value.time_since_epoch()).count(); +} +} + +struct Frame_Sampler::Private { + detail::Gallery_Frame_Atlas atlas; /* 各 Plot 最近原生帧与持久图集的唯一状态源。 */ + std::chrono::nanoseconds period{}; /* 固定采样周期;构造后不变。 */ + double frame_rate_fps{}; /* 对外诊断使用的精确目标频率。 */ + 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 rejected_source_frames{}; /* 无效或过期源帧计数。 */ + std::atomic_uint64_t sampled_frames{}; /* 实际采样计数。 */ + std::atomic_uint64_t early_ticks{}; /* deadline 前被合并的触发计数。 */ + std::atomic_uint64_t missed_periods{}; /* 因迟到直接跳过的完整周期数。 */ + + Private(double target_frame_rate_fps, + std::uint32_t tile_width, + std::uint32_t tile_height, + std::uint32_t columns, + std::vector source_ids, + std::size_t source_slot) + : atlas(tile_width, tile_height, columns, std::move(source_ids)), + period(static_cast(std::llround( + 1'000'000'000.0 / target_frame_rate_fps))), + frame_rate_fps(target_frame_rate_fps), + sampling_source_slot(source_slot) {} +}; + +Frame_Sampler::Frame_Sampler( + double frame_rate_fps, std::uint32_t tile_width, + std::uint32_t tile_height, std::uint32_t columns, + std::vector source_ids, + std::size_t sampling_source_slot) { + if (!std::isfinite(frame_rate_fps) || frame_rate_fps <= 0.0) + throw std::invalid_argument( + "frame sampler rate must be finite and positive"); + if (sampling_source_slot >= source_ids.size()) + throw std::out_of_range("frame sampler source slot is invalid"); + d = std::make_unique( + frame_rate_fps, tile_width, tile_height, columns, + std::move(source_ids), sampling_source_slot); + if (d->period <= std::chrono::nanoseconds::zero()) + throw std::invalid_argument("frame sampler period is too small"); +} + +Frame_Sampler::~Frame_Sampler() = default; + +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); + return Accept_Frame_Result::accepted; + } + d->rejected_source_frames.fetch_add(1, std::memory_order_relaxed); + return result == detail::Gallery_Frame_Atlas::Accept_Frame_Result::stale_frame + ? Accept_Frame_Result::stale_frame + : Accept_Frame_Result::invalid_frame; +} + +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) { + auto composition = d->atlas.compose(); + const auto& source = composition.sources[d->sampling_source_slot]; + Plot_Render_Tick tick{ + now, + source.completion_correlation_id, + static_cast(source.completion_presentation_time.count()) / + 1'000.0, + composition.width, + composition.height}; + return Sampled_Gallery_Frame{ + std::move(composition), tick, deadline_delay_ms}; + }; + auto deadline_ns = d->next_deadline_ns.load(std::memory_order_acquire); + for (;;) { + if (deadline_ns == 0) { + const auto next = now_ns + d->period.count(); + if (!d->next_deadline_ns.compare_exchange_weak( + deadline_ns, next, std::memory_order_acq_rel, + std::memory_order_acquire)) + continue; + d->sampled_frames.fetch_add(1, std::memory_order_relaxed); + return compose_sample(0.0); + } + if (now_ns < deadline_ns) { + d->early_ticks.fetch_add(1, std::memory_order_relaxed); + return std::nullopt; + } + const auto late_ns = now_ns - deadline_ns; + const auto missed = static_cast( + late_ns / d->period.count()); + const auto next = deadline_ns + + static_cast(missed + 1U) * d->period.count(); + if (!d->next_deadline_ns.compare_exchange_weak( + deadline_ns, next, std::memory_order_acq_rel, + std::memory_order_acquire)) + 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); + } +} + +detail::Gallery_Atlas_Description Frame_Sampler::describe() const { + return d->atlas.describe(); +} + +Frame_Sampler_State Frame_Sampler::state() const noexcept { + return { + d->frame_rate_fps, + d->accepted_source_frames.load(std::memory_order_relaxed), + d->rejected_source_frames.load(std::memory_order_relaxed), + d->sampled_frames.load(std::memory_order_relaxed), + d->early_ticks.load(std::memory_order_relaxed), + d->missed_periods.load(std::memory_order_relaxed)}; +} + +} diff --git a/web_server/src/frame_sampling/Frame_Sampler.hpp b/web_server/src/frame_sampling/Frame_Sampler.hpp new file mode 100644 index 0000000..84c7c89 --- /dev/null +++ b/web_server/src/frame_sampling/Frame_Sampler.hpp @@ -0,0 +1,62 @@ +#pragma once +#include "../Plot.hpp" +#include "../detail/Gallery_Frame_Atlas.hpp" +#include +#include +#include +#include +#include +#include +#include + +namespace aethera::web::frame_sampling { + +struct Sampled_Gallery_Frame { + detail::Gallery_Atlas_Composition composition{}; /* 本次采样得到的原生像素图集;下一次成功采样前有效。 */ + Plot_Render_Tick tick{}; /* 触发采样的服务端帧时钟。 */ + double deadline_delay_ms{}; /* 实际采样相对固定周期 deadline 的延迟。 */ +}; + +struct Frame_Sampler_State { + double target_frame_rate_fps{}; /* 采样器固定目标频率。 */ + std::uint64_t accepted_source_frames{}; /* 已进入最近帧快照的 Plot 完成帧数。 */ + std::uint64_t rejected_source_frames{}; /* 因尺寸、布局或顺序不合法而拒绝的完成帧数。 */ + std::uint64_t sampled_frames{}; /* 到达 deadline 并实际生成的图集样本数。 */ + std::uint64_t early_ticks{}; /* deadline 前到达、未触发采样的外接图次数。 */ + std::uint64_t missed_periods{}; /* 调度延迟跨过的完整采样周期数。 */ +}; + +/* + * Gallery 像素采样的唯一权威入口。Plot 可以按任意帧策略发布,采样器只在 + * 固定 deadline 到达时读取每个槽位的 latest,不排队、不补历史帧。 + */ +struct Frame_Sampler final { +public: + enum struct Accept_Frame_Result : std::uint8_t { + accepted, + invalid_frame, + stale_frame + }; + + Frame_Sampler(double frame_rate_fps, + std::uint32_t tile_width, + std::uint32_t tile_height, + std::uint32_t columns, + std::vector source_ids, + std::size_t sampling_source_slot); + ~Frame_Sampler(); + Frame_Sampler(const Frame_Sampler&) = delete; + Frame_Sampler& operator=(const Frame_Sampler&) = delete; + + [[nodiscard]] Accept_Frame_Result accept_frame( + std::size_t slot, std::shared_ptr frame); + [[nodiscard]] std::optional sample(); + [[nodiscard]] detail::Gallery_Atlas_Description describe() const; + [[nodiscard]] Frame_Sampler_State state() const noexcept; + +private: + struct Private; + std::unique_ptr d; +}; + +} diff --git a/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.cpp b/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.cpp new file mode 100644 index 0000000..f06d4b7 --- /dev/null +++ b/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.cpp @@ -0,0 +1,39 @@ +#include "FFmpeg_Frame_Transport.hpp" +#include +#include + +namespace aethera::web::frame_sampling::ffmpeg { + +struct FFmpeg_Frame_Transport::Private { + H264_Encoder encoder; /* FFmpeg 编码上下文的唯一所有者。 */ + std::atomic_bool key_frame_requested{}; /* 下一次实际编码样本是否强制关键帧。 */ + + explicit Private(double frame_rate_fps) : encoder(frame_rate_fps) {} +}; + +FFmpeg_Frame_Transport::FFmpeg_Frame_Transport(double frame_rate_fps) + : d(std::make_unique(frame_rate_fps)) {} + +FFmpeg_Frame_Transport::~FFmpeg_Frame_Transport() = default; + +void FFmpeg_Frame_Transport::request_key_frame() noexcept { + d->key_frame_requested.store(true, std::memory_order_release); +} + +std::optional FFmpeg_Frame_Transport::encode( + const Sampled_Gallery_Frame& frame) { + if (d->key_frame_requested.exchange(false, std::memory_order_acq_rel)) + d->encoder.request_key_frame(); + return d->encoder.encode( + frame.composition.pixels, + frame.composition.width, + frame.composition.height, + frame.composition.layout == Plot_Pixel_Layout::bgra8 + ? Video_Pixel_Layout::bgra + : Video_Pixel_Layout::rgba, + frame.tick.sequence, + std::chrono::microseconds{static_cast(std::llround( + frame.tick.time_milliseconds * 1'000.0))}); +} + +} diff --git a/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.hpp b/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.hpp new file mode 100644 index 0000000..e0557e0 --- /dev/null +++ b/web_server/src/frame_sampling/ffmpeg/FFmpeg_Frame_Transport.hpp @@ -0,0 +1,26 @@ +#pragma once +#include "../../H264_Encoder.hpp" +#include "../Frame_Sampler.hpp" +#include +#include + +namespace aethera::web::frame_sampling::ffmpeg { + +/* 固定采样帧到 FFmpeg H.264 access unit 的单一编码位置。 */ +struct FFmpeg_Frame_Transport final { +public: + explicit FFmpeg_Frame_Transport(double frame_rate_fps); + ~FFmpeg_Frame_Transport(); + FFmpeg_Frame_Transport(const FFmpeg_Frame_Transport&) = delete; + FFmpeg_Frame_Transport& operator=(const FFmpeg_Frame_Transport&) = delete; + + void request_key_frame() noexcept; + [[nodiscard]] std::optional encode( + const Sampled_Gallery_Frame& frame); + +private: + struct Private; + std::unique_ptr d; +}; + +} diff --git a/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.cpp b/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.cpp new file mode 100644 index 0000000..c43baed --- /dev/null +++ b/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.cpp @@ -0,0 +1,68 @@ +#include "WebSocket_Pixel_Frame.hpp" +#include +#include +#include + +namespace aethera::web::frame_sampling::websocket_pixel { +namespace { +constexpr std::size_t header_size{48}; + +template +void append_little_endian(std::string& output, Integer value) { + using Unsigned = std::make_unsigned_t; + auto remaining = static_cast(value); + for (std::size_t index = 0; index < sizeof(Integer); ++index) { + output.push_back(static_cast(remaining & 0xffU)); + remaining >>= 8U; + } +} +} + +WebSocket_Pixel_Frame pack_websocket_pixel_frame( + const Sampled_Gallery_Frame& frame) { + const auto& composition = frame.composition; + const auto expected = static_cast(composition.width) * + composition.height * 4U; + if (composition.width == 0 || composition.height == 0 || + composition.pixels.size() != expected) + throw std::invalid_argument( + "WebSocket pixel frame requires a complete four-byte atlas"); + if (expected > std::numeric_limits::max()) + throw std::overflow_error("WebSocket pixel payload exceeds protocol size"); + + const auto format = composition.layout == Plot_Pixel_Layout::bgra8 + ? Wire_Pixel_Format::bgra8_premultiplied + : Wire_Pixel_Format::rgba8; + WebSocket_Pixel_Frame result; + result.sequence = frame.tick.sequence; + result.format = format; + result.bytes.reserve(header_size + expected); + result.bytes.append("AERP", 4); + append_little_endian(result.bytes, 1U); + append_little_endian( + result.bytes, static_cast(header_size)); + append_little_endian( + result.bytes, static_cast(format)); + append_little_endian(result.bytes, 0U); + append_little_endian(result.bytes, 0U); + append_little_endian(result.bytes, composition.width); + append_little_endian(result.bytes, composition.height); + append_little_endian( + result.bytes, composition.width * 4U); + append_little_endian( + result.bytes, static_cast(expected)); + append_little_endian(result.bytes, frame.tick.sequence); + append_little_endian(result.bytes, + static_cast(std::max(0, + std::chrono::duration_cast( + std::chrono::duration( + frame.tick.time_milliseconds)).count()))); + append_little_endian(result.bytes, 0U); + if (result.bytes.size() != header_size) + throw std::logic_error("WebSocket pixel protocol header size mismatch"); + result.bytes.append( + reinterpret_cast(composition.pixels.data()), expected); + return result; +} + +} diff --git a/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.hpp b/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.hpp new file mode 100644 index 0000000..197a33c --- /dev/null +++ b/web_server/src/frame_sampling/websocket_pixel/WebSocket_Pixel_Frame.hpp @@ -0,0 +1,24 @@ +#pragma once +#include "../Frame_Sampler.hpp" +#include +#include +#include + +namespace aethera::web::frame_sampling::websocket_pixel { + +enum struct Wire_Pixel_Format : std::uint8_t { + rgba8 = 0, + bgra8_premultiplied = 1 +}; + +struct WebSocket_Pixel_Frame { + std::string bytes{}; /* 48 字节小端头部及未转换的原生像素负载。 */ + std::uint64_t sequence{}; /* 本二进制包对应的采样时钟序号。 */ + Wire_Pixel_Format format{Wire_Pixel_Format::rgba8}; /* 浏览器解释负载通道与 alpha 的依据。 */ +}; + +/* 只封装协议头并复制所有权,不做缩放、通道交换或 alpha 转换。 */ +[[nodiscard]] WebSocket_Pixel_Frame pack_websocket_pixel_frame( + const Sampled_Gallery_Frame& frame); + +} diff --git a/web_server/tests/Frame_Sampler_Tests.cpp b/web_server/tests/Frame_Sampler_Tests.cpp new file mode 100644 index 0000000..85b4443 --- /dev/null +++ b/web_server/tests/Frame_Sampler_Tests.cpp @@ -0,0 +1,59 @@ +#include +#include +#include +#include +#include +#include +#include + +namespace aethera::web::frame_sampling { +namespace { +std::shared_ptr native_bgra_frame() { + auto pixels = std::make_shared>(16U); + (*pixels)[0] = std::byte{11}; + (*pixels)[1] = std::byte{22}; + (*pixels)[2] = std::byte{33}; + (*pixels)[3] = std::byte{44}; + return std::make_shared(Plot_Pixel_Frame{ + std::move(pixels), Plot_Pixel_Layout::bgra8, + std::chrono::microseconds{12'345}, 1, 9, 1, 9, 2, 2}); +} +} + +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); + const auto sample = sampler.sample(); + ASSERT_TRUE(sample); + 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}); + EXPECT_EQ(sample->composition.pixels[2], std::byte{33}); + EXPECT_EQ(sample->composition.pixels[3], std::byte{44}); + EXPECT_EQ(sample->tick.sequence, 9U); + EXPECT_DOUBLE_EQ(sample->tick.time_milliseconds, 12.345); + EXPECT_FALSE(sampler.sample()); + 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); +} + +TEST(WebSocket_Pixel_Frame, Header_Describes_Unchanged_Native_Payload) { + 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); + const auto sample = sampler.sample(); + ASSERT_TRUE(sample); + const auto packet = websocket_pixel::pack_websocket_pixel_frame(*sample); + ASSERT_EQ(packet.bytes.size(), 48U + 16U); + EXPECT_EQ(packet.bytes.substr(0, 4), "AERP"); + EXPECT_EQ(static_cast(packet.bytes[8]), 1U); + EXPECT_EQ(static_cast(packet.bytes[48]), 11U); + EXPECT_EQ(static_cast(packet.bytes[49]), 22U); + EXPECT_EQ(static_cast(packet.bytes[50]), 33U); + EXPECT_EQ(static_cast(packet.bytes[51]), 44U); +} + +} diff --git a/web_server/tests/Gallery_Frame_Clock_Tests.cpp b/web_server/tests/Gallery_Frame_Clock_Tests.cpp deleted file mode 100644 index 08d94a9..0000000 --- a/web_server/tests/Gallery_Frame_Clock_Tests.cpp +++ /dev/null @@ -1,61 +0,0 @@ -#include -#include -#include -#include -#include -#include -#include -#include - -namespace aethera::web::detail { -TEST(Gallery_Frame_Clock, Produces_One_Monotonic_Absolute_Timeline) { - std::mutex mutex; - std::condition_variable condition; - std::vector ticks; - std::exception_ptr failure; - { - Gallery_Frame_Clock clock( - 100.0, - [&](Plot_Render_Tick tick) { - std::lock_guard lock(mutex); - ticks.push_back(tick); - condition.notify_all(); - }, - [&](std::exception_ptr value) { - std::lock_guard lock(mutex); - failure = std::move(value); - condition.notify_all(); - }); - clock.start(); - clock.start(); - { - std::unique_lock lock(mutex); - ASSERT_TRUE(condition.wait_for(lock, std::chrono::seconds(2), [&] { - return ticks.size() >= 8 || failure; - })); - ASSERT_FALSE(failure); - ASSERT_GE(ticks.size(), 8U); - for (std::size_t index = 1; index < ticks.size(); ++index) { - EXPECT_GT(ticks[index].sequence, ticks[index - 1].sequence); - EXPECT_GT(ticks[index].time_milliseconds, - ticks[index - 1].time_milliseconds); - } - EXPECT_NEAR(ticks.front().time_milliseconds, - static_cast(ticks.front().sequence) * 10.0, - 0.001); - } - clock.stop(); - std::this_thread::sleep_for(std::chrono::milliseconds(30)); - std::size_t stopped_count{}; - { - std::lock_guard lock(mutex); - stopped_count = ticks.size(); - } - std::this_thread::sleep_for(std::chrono::milliseconds(30)); - { - std::lock_guard lock(mutex); - EXPECT_EQ(ticks.size(), stopped_count); - } - } -} -} diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index f5a53d2..1f3ce0d 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -30,6 +30,7 @@ type Plot_Execution_Policies = Record; type Frame_Analysis = Omit & {data_generator?: Data_Generator}; type Schema = {protocol: "aethera.plot.inspector"; version: 2; components: Component[]; frame_analysis: Frame_Analysis}; type Stream_Status = "IDLE" | "CONNECTING" | "LIVE" | "OFFLINE"; +type Gallery_Transport_Mode = "ffmpeg" | "websocket_pixels"; type Frame_Pacing_Mode = "manual" | "fixed_rate" | "maximum_rate"; type Frame_Delivery = "gallery-video" | "diagnostics"; type Pixel_Format = "rgba8" | "bgra8_premultiplied"; @@ -53,18 +54,23 @@ type Video_Playback_Metrics = {frame_rate_fps: number; presented_frames: number; jitter_buffer_ms: number; decode_processing_ms: number; estimated_playout_delay_ms: number; freeze_count: number; current_time_seconds: number; ready_state: number}; type Gallery_Tile = {column: number; row: number}; -type Gallery_Layout = {kind: "gallery_layout"; protocol: "aethera.gallery.video"; version: 2; columns: number; rows: number; +type Gallery_Layout = {kind: "gallery_layout"; protocol: "aethera.gallery.video"; version: 2 | 3; + transport?: Gallery_Transport_Mode; frame_rate_fps?: number; columns: number; rows: number; 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}; -type Gallery_Transport_Metrics = {kind: "gallery_metrics"; protocol: "aethera.gallery.video"; version: 2; clock_sequence: number; - encoded_frame_count: number; encoder_backend: "nvenc" | "vulkan_video"; target_frame_rate_fps: number; clock_delivery_rate_fps: number; encoded_frame_rate_fps: 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; + encoded_frame_rate_fps: number; pixel_frame_rate_fps?: number; compose_average_ms: number; compose_p95_ms: number; encode_average_ms: number; encode_p95_ms: number; - publish_average_ms: number; publish_p95_ms: number; encoded_bytes: number; fresh_tiles: number; missing_tiles: number; - rejected_frames: number; skipped_encode_ticks: number; sources: Record}; + pixel_pack_average_ms?: number; pixel_pack_p95_ms?: number; sample_delay_average_ms?: number; sample_delay_p95_ms?: number; + publish_average_ms: number; publish_p95_ms: number; encoded_bytes: number; pixel_bytes?: number; fresh_tiles: number; missing_tiles: number; + rejected_frames: number; skipped_encode_ticks?: number; skipped_sample_ticks?: number; sources: Record}; type Gallery_Video_State = {status: Stream_Status; stream: MediaStream | null; layout: Gallery_Layout | null; - playback: Video_Playback_Metrics; transport: Gallery_Transport_Metrics | null; error: string | null}; + pixel_surface: Gallery_Pixel_Surface | null; playback: Video_Playback_Metrics; + transport: Gallery_Transport_Metrics | null; error: string | null}; type Gallery_Video_States = Record; type Frame_Diagnostics = {server: Plot_Diagnostics; samples: Frame_Sample[]; video_playback: Video_Playback_Metrics; browser_input_statistics: Browser_Input_Statistics}; @@ -85,7 +91,8 @@ type Taskflow_Execution_Trace = {native_id: string; node_id: string; worker_id: type Taskflow_Frame_Trace = {sequence: number; correlation_id: number; created_time_unix_ns: number; worker_count: number; markers?: Record; graphs: Taskflow_Graph_Trace[]; executions: Taskflow_Execution_Trace[]}; type Taskflow_Frame_Response = {protocol: "aethera.taskflow.frames"; version: 1; requested: number; remaining: number; - captured: number; complete: boolean; frames: Taskflow_Frame_Trace[]}; + 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; peak_queue_size: number; max_queue_capacity: number; active_task: {native_id: string; type: string; time_ns: number}; @@ -110,9 +117,10 @@ function socket_url(path: string) { return url.toString(); } const gallery_connection_id = crypto.randomUUID(); -function gallery_socket_url(path: string) { +function gallery_socket_url(path: string, transport: Gallery_Transport_Mode) { const url = new URL(socket_url(path)); url.searchParams.set("connection", gallery_connection_id); + url.searchParams.set("transport", transport); return url.toString(); } @@ -128,7 +136,7 @@ function valid_gallery_layout(value: unknown): value is Gallery_Layout { if (!value || typeof value !== "object") return false; const layout = value as Partial; return layout.kind === "gallery_layout" && - layout.protocol === "aethera.gallery.video" && layout.version === 2 && + layout.protocol === "aethera.gallery.video" && (layout.version === 2 || layout.version === 3) && typeof layout.tile_width === "number" && typeof layout.tile_height === "number" && typeof layout.width === "number" && typeof layout.height === "number" && Boolean(layout.plots); @@ -138,7 +146,7 @@ function valid_gallery_metrics(value: unknown): value is Gallery_Transport_Metri if (!value || typeof value !== "object") return false; const metrics = value as Partial; return metrics.kind === "gallery_metrics" && - metrics.protocol === "aethera.gallery.video" && metrics.version === 2 && + metrics.protocol === "aethera.gallery.video" && (metrics.version === 2 || metrics.version === 3) && typeof metrics.encoded_frame_rate_fps === "number" && typeof metrics.clock_delivery_rate_fps === "number" && Boolean(metrics.sources); @@ -194,12 +202,157 @@ const empty_video_playback = (): Video_Playback_Metrics => ({ current_time_seconds: 0, ready_state: 0 }); +type Gallery_Pixel_Packet = {format: number; width: number; height: number; stride: number; + payload_bytes: number; sequence: number; presentation_microseconds: number; pixels: Uint8Array}; + +function parse_gallery_pixel_packet(buffer: ArrayBuffer): Gallery_Pixel_Packet { + if (buffer.byteLength < 48) throw new Error("原始像素帧头部不完整"); + const bytes = new Uint8Array(buffer); + if (bytes[0] !== 0x41 || bytes[1] !== 0x45 || bytes[2] !== 0x52 || bytes[3] !== 0x50) + throw new Error("原始像素帧 magic 不匹配"); + const view = new DataView(buffer); + const version = view.getUint16(4, true); + const header_bytes = view.getUint16(6, true); + const format = view.getUint8(8); + const width = view.getUint32(12, true); + const height = view.getUint32(16, true); + const stride = view.getUint32(20, true); + const payload_bytes = view.getUint32(24, true); + if (version !== 1 || header_bytes !== 48 || (format !== 0 && format !== 1) || + width === 0 || height === 0 || stride !== width * 4 || + payload_bytes !== stride * height || buffer.byteLength !== header_bytes + payload_bytes) + throw new Error("原始像素帧协议字段无效"); + return {format, width, height, stride, payload_bytes, + sequence: Number(view.getBigUint64(28, true)), + presentation_microseconds: Number(view.getBigUint64(36, true)), + pixels: new Uint8Array(buffer, header_bytes, payload_bytes)}; +} + +function compile_gallery_shader(gl: WebGL2RenderingContext, type: number, source: string) { + const shader = gl.createShader(type); + if (!shader) throw new Error("无法创建原始像素 WebGL shader"); + gl.shaderSource(shader, source); + gl.compileShader(shader); + if (!gl.getShaderParameter(shader, gl.COMPILE_STATUS)) { + const reason = gl.getShaderInfoLog(shader) ?? "unknown shader error"; + gl.deleteShader(shader); + throw new Error(`原始像素 shader 编译失败: ${reason}`); + } + return shader; +} + +/* 每个媒体组只上传一次图集纹理;各 Plot Canvas 只裁剪自己的 tile。 */ +class Gallery_Pixel_Surface { + readonly canvas = document.createElement("canvas"); + private readonly gl: WebGL2RenderingContext; + private readonly texture: WebGLTexture; + private readonly format_location: WebGLUniformLocation; + private texture_width = 0; + private texture_height = 0; + private presented_frames = 0; + private previous_presented_frames = 0; + private previous_sample_time = performance.now(); + private latest_presentation_microseconds = 0; + + constructor() { + const gl = this.canvas.getContext("webgl2", { + alpha: true, premultipliedAlpha: true, preserveDrawingBuffer: true, + antialias: false, depth: false, stencil: false + }); + if (!gl) throw new Error("WebSocket 原始像素模式需要 WebGL2"); + this.gl = gl; + const vertex = compile_gallery_shader(gl, gl.VERTEX_SHADER, `#version 300 es + in vec2 position; in vec2 texture_coordinate; out vec2 uv; + void main() { gl_Position = vec4(position, 0.0, 1.0); uv = texture_coordinate; }`); + const fragment = compile_gallery_shader(gl, gl.FRAGMENT_SHADER, `#version 300 es + precision mediump float; in vec2 uv; uniform sampler2D pixels; + uniform int pixel_format; out vec4 color; + void main() { vec4 value = texture(pixels, uv); + color = pixel_format == 1 ? value.bgra : value; }`); + const program = gl.createProgram(); + if (!program) throw new Error("无法创建原始像素 WebGL program"); + gl.attachShader(program, vertex); gl.attachShader(program, fragment); gl.linkProgram(program); + gl.deleteShader(vertex); gl.deleteShader(fragment); + if (!gl.getProgramParameter(program, gl.LINK_STATUS)) + throw new Error(`原始像素 program 链接失败: ${gl.getProgramInfoLog(program) ?? "unknown"}`); + gl.useProgram(program); + const position = gl.getAttribLocation(program, "position"); + const texture_coordinate = gl.getAttribLocation(program, "texture_coordinate"); + const buffer = gl.createBuffer(); + if (!buffer) throw new Error("无法创建原始像素顶点缓冲"); + gl.bindBuffer(gl.ARRAY_BUFFER, buffer); + gl.bufferData(gl.ARRAY_BUFFER, new Float32Array([ + -1, -1, 0, 1, 1, -1, 1, 1, -1, 1, 0, 0, 1, 1, 1, 0 + ]), gl.STATIC_DRAW); + gl.enableVertexAttribArray(position); gl.vertexAttribPointer(position, 2, gl.FLOAT, false, 16, 0); + gl.enableVertexAttribArray(texture_coordinate); + gl.vertexAttribPointer(texture_coordinate, 2, gl.FLOAT, false, 16, 8); + const texture = gl.createTexture(); + if (!texture) throw new Error("无法创建原始像素纹理"); + this.texture = texture; + gl.bindTexture(gl.TEXTURE_2D, texture); + gl.texParameteri(gl.TEXTURE_2D, gl.TEXTURE_MIN_FILTER, gl.LINEAR); + gl.texParameteri(gl.TEXTURE_2D, gl.TEXTURE_MAG_FILTER, gl.LINEAR); + gl.texParameteri(gl.TEXTURE_2D, gl.TEXTURE_WRAP_S, gl.CLAMP_TO_EDGE); + gl.texParameteri(gl.TEXTURE_2D, gl.TEXTURE_WRAP_T, gl.CLAMP_TO_EDGE); + gl.uniform1i(gl.getUniformLocation(program, "pixels"), 0); + const format_location = gl.getUniformLocation(program, "pixel_format"); + if (!format_location) throw new Error("原始像素 shader 缺少 pixel_format"); + this.format_location = format_location; + } + + render(buffer: ArrayBuffer) { + const packet = parse_gallery_pixel_packet(buffer); + const gl = this.gl; + if (this.texture_width !== packet.width || this.texture_height !== packet.height) { + this.canvas.width = packet.width; this.canvas.height = packet.height; + gl.viewport(0, 0, packet.width, packet.height); + this.texture_width = packet.width; this.texture_height = packet.height; + gl.texImage2D(gl.TEXTURE_2D, 0, gl.RGBA8, packet.width, packet.height, + 0, gl.RGBA, gl.UNSIGNED_BYTE, packet.pixels); + } else { + gl.texSubImage2D(gl.TEXTURE_2D, 0, 0, 0, packet.width, packet.height, + gl.RGBA, gl.UNSIGNED_BYTE, packet.pixels); + } + gl.uniform1i(this.format_location, packet.format); + gl.drawArrays(gl.TRIANGLE_STRIP, 0, 4); + this.presented_frames += 1; + this.latest_presentation_microseconds = packet.presentation_microseconds; + return packet.sequence; + } + + draw_tile(target: HTMLCanvasElement, tile: Gallery_Tile, width: number, height: number) { + if (target.width !== width) target.width = width; + if (target.height !== height) target.height = height; + const context = target.getContext("2d", {alpha: true}); + if (!context) throw new Error("无法创建 Plot Canvas 2D context"); + context.clearRect(0, 0, width, height); + context.drawImage(this.canvas, tile.column * width, tile.row * height, + width, height, 0, 0, width, height); + } + + playback(): Video_Playback_Metrics { + const now = performance.now(); + const elapsed = now - this.previous_sample_time; + const frame_rate_fps = elapsed > 0 + ? (this.presented_frames - this.previous_presented_frames) * 1000 / elapsed : 0; + this.previous_sample_time = now; + this.previous_presented_frames = this.presented_frames; + return {...empty_video_playback(), frame_rate_fps, + presented_frames: this.presented_frames, + current_time_seconds: this.latest_presentation_microseconds / 1_000_000, + ready_state: this.presented_frames > 0 ? 4 : 0}; + } + + dispose() { this.gl.deleteTexture(this.texture); } +} + const connecting_gallery_video = (): Gallery_Video_State => ({ - status: "CONNECTING", stream: null, layout: null, + status: "CONNECTING", stream: null, layout: null, pixel_surface: null, playback: empty_video_playback(), transport: null, error: null }); -function use_gallery_videos(plots: Plot[]): Gallery_Video_States { +function use_gallery_videos(plots: Plot[], transport_mode: Gallery_Transport_Mode): Gallery_Video_States { const media = useMemo(() => [...new Set(plots.map(plot => plot.media))].sort(), [plots]); const media_key = media.join("\n"); const [states, set_states] = useState({}); @@ -215,6 +368,9 @@ function use_gallery_videos(plots: Plot[]): Gallery_Video_States { let stopped = false; let peer: RTCPeerConnection | null = null; let signal_chain = Promise.resolve(); + let pixel_chain = Promise.resolve(); + let pixel_surface: Gallery_Pixel_Surface | null = null; + let pixel_live = false; let previous_time = performance.now(); let previous_decoded = 0; const update = (change: (current: Gallery_Video_State) => Gallery_Video_State) => @@ -232,11 +388,19 @@ function use_gallery_videos(plots: Plot[]): Gallery_Video_States { }; window.addEventListener("aethera:gallery-diagnostics", receive_server_diagnostics); - const socket = new ReconnectingWebSocket(gallery_socket_url(endpoint), [], { + const socket = new ReconnectingWebSocket( + gallery_socket_url(endpoint, transport_mode), [], { minReconnectionDelay: 300, maxReconnectionDelay: 5000, reconnectionDelayGrowFactor: 1.6, maxRetries: Number.POSITIVE_INFINITY }); + socket.binaryType = "arraybuffer"; const sample_stats = async () => { + if (transport_mode === "websocket_pixels") { + if (pixel_surface) + update(current => ({...current, + playback: pixel_surface!.playback()})); + return; + } if (!peer) return; try { const reports = await peer.getStats(); @@ -313,16 +477,42 @@ function use_gallery_videos(plots: Plot[]): Gallery_Video_States { }; }; socket.onopen = () => { - create_peer(); - update(current => ({...current, status: "CONNECTING", error: null})); + if (transport_mode === "ffmpeg") create_peer(); + else { + pixel_surface?.dispose(); + pixel_surface = new Gallery_Pixel_Surface(); + pixel_live = false; + } + update(current => ({...current, status: "CONNECTING", error: null, + stream: null, pixel_surface})); }; socket.onclose = () => { peer?.close(); peer = null; - if (!stopped) update(current => ({...current, status: "CONNECTING", stream: null})); + if (!stopped) update(current => ({...current, + status: "CONNECTING", stream: null})); }; socket.onmessage = event => { - if (typeof event.data !== "string") return; + if (typeof event.data !== "string") { + if (transport_mode !== "websocket_pixels" || !pixel_surface) return; + pixel_chain = pixel_chain.then(async () => { + const buffer = event.data instanceof ArrayBuffer + ? event.data : event.data instanceof Blob + ? await event.data.arrayBuffer() : null; + if (!buffer || stopped || !pixel_surface) return; + const sequence = pixel_surface.render(buffer); + if (!pixel_live) { + pixel_live = true; + update(current => ({...current, status: "LIVE", + pixel_surface, error: null})); + } + window.dispatchEvent(new CustomEvent("aethera:gallery-pixel-frame", + {detail: {media: endpoint, surface: pixel_surface, sequence}})); + }).catch(error => update(current => ({...current, + status: "OFFLINE", error: error instanceof Error + ? error.message : "原始像素帧处理失败"}))); + return; + } let decoded: unknown; try { decoded = JSON.parse(event.data); } catch { return; } if (!decoded || typeof decoded !== "object") return; @@ -363,11 +553,12 @@ function use_gallery_videos(plots: Plot[]): Gallery_Video_States { window.removeEventListener("aethera:gallery-diagnostics", receive_server_diagnostics); peer?.close(); + pixel_surface?.dispose(); socket.close(); }; }); return () => cleanups.forEach(cleanup => cleanup()); - }, [media_key]); + }, [media_key, transport_mode]); return states; } @@ -1110,7 +1301,7 @@ function Frame_Policy_Pane({plot, analysis, busy, on_refresh, on_update, on_manu }) { const fields = analysis?.fields.filter(field => field.editable) ?? []; return
-
{plot.dimension} 后端帧策略渲染、采样、H.264 编码与 WebRTC 传输全部由当前 Plot 的后端策略控制;浏览器隐藏画面不会停采样或断流。默认全部开启,用于完整链路压测。
+
{plot.dimension} 后端帧策略Scene 按当前后端策略渲染;媒体层独立按 30 FPS 读取 latest,并选择 FFmpeg/WebRTC 或 WebSocket 原始像素。浏览器隐藏画面不会停采样。
{analysis ?
{fields.map(field => on_update(analysis, field, value)}/>)}
:

正在读取帧策略…

}
; @@ -1174,9 +1365,12 @@ function taskflow_node_name(name: string) { if (name === "scene.paint.complete") return "Scene · 像素绘制完成"; if (name === "scene.completion") return "Scene · 帧完成处理"; if (name === "plot.frame.publish") return "Plot · 发布完成帧"; - if (name === "gallery.atlas.compose") return "Gallery · 图集合成"; + if (name === "gallery.sample.capture") return "Gallery · 30 FPS 最近帧采样"; if (name === "FFmpeg.H264.encode" || name === "gallery.h264.encode") return "FFmpeg H.264 · 编码"; - if (name === "gallery.webrtc.publish") return "WebRTC · 视频帧发布"; + if (name === "gallery.ffmpeg.webrtc.publish") return "WebRTC · H.264 帧发布"; + if (name === "gallery.websocket_pixels.pack") return "WebSocket · 原生像素封装"; + if (name === "gallery.websocket_pixels.publish") return "WebSocket · 原生像素发布"; + if (name === "gallery.sample.complete") return "Gallery · 采样完成"; const owner = parts[0]; const stage = parts[1]; @@ -1630,10 +1824,15 @@ function Taskflow_Dag({graph, executions, frame, aggregate, components, gallery_ correlation_id: frame?.correlation_id, markers: frame?.markers ?? {}, graph, + /* 保留当前选中拓扑,同时把本帧其余真实 Task_Graph 一并复制。 + * 这样默认停在 render_2d.frame 时,诊断 JSON 里仍能直接检索 + * gallery.media / FFmpeg.H264.encode,不会误以为媒体拓扑丢失。 */ + graphs: frame?.graphs ?? [graph], aggregate: aggregate ? {frames: aggregate.frames, samples: aggregate.samples, topology_variants: aggregate.topology_variants, wall: aggregate.wall, nodes: Object.fromEntries(aggregate.nodes)} : undefined, - executions: executions.filter(value => native_ids.has(value.native_id)) + executions: executions.filter(value => native_ids.has(value.native_id)), + all_executions: frame?.executions ?? executions }, null, 2)); set_copy_state("已复制"); window.setTimeout(() => set_copy_state("复制拓扑 JSON"), 1200); @@ -1758,8 +1957,11 @@ function Taskflow_Frame_Pane({plot, components}: {plot: Plot; components: Compon
- - {response ? `${response.captured}/${response.requested} 帧${response.complete ? " · 已完成" : ` · 还需 ${response.remaining} 帧`}` : "按需捕获,未请求时 Observer 不写入逐帧数据"}
+ + {response ? `${response.captured}/${response.requested} Plot 帧${response.media_requested !== undefined + ? ` · 媒体 ${response.media_captured ?? 0}/${response.media_requested}` : ""}${response.complete + ? " · 已完成" : ` · 还需 ${response.remaining} 帧`}` + : "捕获当前 Plot 全部拓扑,并附加 Gallery 实际执行的媒体 Task_Graph"}
{error ?

{error}

: null} {frame ? <>
diff --git a/webapp_gallery/src/styles.css b/webapp_gallery/src/styles.css index 5d4daa0..dfef2f4 100644 --- a/webapp_gallery/src/styles.css +++ b/webapp_gallery/src/styles.css @@ -36,6 +36,8 @@ h2 { margin: 4px 0 0; font-size: 23px; letter-spacing: -.03em; } .error { color: #ff879d !important; } .topbarActions { display: flex; align-items: center; justify-content: flex-end; gap: 8px; } +.transportSelector { display: flex; align-items: center; gap: 7px; color: #8296b2; font-size: 11px; white-space: nowrap; } +.transportSelector select { padding: 7px 9px; color: #c9d8eb; border: 1px solid #2a3b55; border-radius: 8px; background: #0a1220; } .topbarActions button, nav button { padding: 8px 13px; color: #aabbd3; @@ -78,7 +80,7 @@ nav { display: flex; flex-wrap: wrap; gap: 8px; padding: 18px 0; } .plotExecutionPolicy label { display:flex; align-items:center; gap:5px; cursor:pointer; user-select:none; } .plotExecutionPolicy input { accent-color:#52dbc1; } .plotViewportHidden { position:relative; } -.plotViewportHidden video { visibility:hidden; } +.plotViewportHidden video, .plotViewportHidden canvas { visibility:hidden; } .plotHiddenState { position:absolute; inset:0; display:grid; place-content:center; gap:5px; text-align:center; color:#91a5bf; background:repeating-linear-gradient(135deg,rgba(15,27,43,.96),rgba(15,27,43,.96) 12px,rgba(18,33,52,.96) 12px,rgba(18,33,52,.96) 24px); } .plotHiddenState strong { color:#d9e8f6; } canvas { display: block; width: 100%; height: 100%; background: #070d18; } diff --git a/绘制优化.md b/绘制优化.md new file mode 100644 index 0000000..d7de1ff --- /dev/null +++ b/绘制优化.md @@ -0,0 +1,386 @@ +可以实现,而且你这个“天然分区并行”的场景很适合这么做。但要把边界说清楚:这不是 Blend2D 官方提供的 `BLImageView` 功能,而是利用 +`BLImage::create_from_data()` 给同一块外部 framebuffer 建多个独立的 `BLImage` 包装。Blend2D 官方明确支持外部 +`pixel_data + stride` 创建 image。 ([Blend2D][1]) + +比如最终 framebuffer 是 1920×1080 PRGB32: + +```cpp +struct Framebuffer { + uint8_t* data; + int width; + int height; + intptr_t stride; +}; + +BLImage make_view( + Framebuffer& fb, + int x, + int y, + int width, + int height) +{ + BLImage image; + uint8_t* ptr = fb.data + y * fb.stride + x * 4; + + image.create_from_data( + width, + height, + BL_FORMAT_PRGB32, + ptr, + fb.stride, + BL_DATA_ACCESS_RW); + + return image; +} +``` + +比如: + +```text +1920 × 1080 framebuffer + +┌──────────────┬──────────────┐ +│ view A │ view B │ +│ 960×540 │ 960×540 │ +├──────────────┼──────────────┤ +│ view C │ view D │ +│ 960×540 │ 960×540 │ +└──────────────┴──────────────┘ +``` + +然后: + +```cpp +auto a = make_view(fb, 0, 0, 960, 540); +auto b = make_view(fb, 960, 0, 960, 540); +auto c = make_view(fb, 0, 540, 960, 540); +auto d = make_view(fb, 960, 540, 960, 540); +``` + +Taskflow: + +```text +worker 0 → BLContext(a) +worker 1 → BLContext(b) +worker 2 → BLContext(c) +worker 3 → BLContext(d) +``` + +四个 context 都可以是同步单线程 context: + +```cpp +BLContext ctx(a); +``` + +Blend2D 文档说明这种普通 `BLContext(image)` 就是 single-threaded synchronous rendering。 ([Blend2D][2]) + +最关键的是为什么它可以绕开“一个 image 只能有一个 renderer”。 + +因为 Blend2D 的 exclusive-access 管理针对的是它自己的 image data 对象。`BLContext::begin()` 会增加目标 image 的 writer +count,并保证同一个 image data 不被另一个 renderer 同时使用。 ([Blend2D][2]) + +而你这里创建的是: + +```text +BLImage A → external-memory wrapper A +BLImage B → external-memory wrapper B +BLImage C → external-memory wrapper C +``` + +它们不是: + +```cpp +BLImage b = a; +``` + +这种 weak copy。 + +每次都是独立调用: + +```cpp +create_from_data(...) +``` + +所以 Blend2D 看到的是三个不同 image data 对象。 + +只是: + +```text +物理内存: + +A.data ─┐ +B.data ─┼→ 同一个 framebuffer allocation +C.data ─┘ +``` + +Blend2D并不知道这些外部内存背后属于同一 allocation。 + +因此真正的线程安全责任就在你这里。 + +你必须保证: + +```text +A 实际写入地址集合 +∩ +B 实际写入地址集合 += +∅ +``` + +如果做到这一点,从 C++ 内存访问角度就是正常的并行写不同内存。 + +你需要特别注意几个问题。 + +第一,`stride` 必须是父 framebuffer 的 stride。 + +例如: + +```text +parent width = 1920 +PRGB32 = 4 bytes +stride = 7680 +``` + +B 虽然只有 960 宽: + +```cpp +pixel_data = base + 960 * 4; +``` + +但: + +```cpp +stride = 1920 * 4; +``` + +仍然不能写成: + +```cpp +stride = 960 * 4; +``` + +否则下一行就跑错位置。 + +`BLImageData` 本身也是这么定义的:`pixel_data` 是图像左上角,`stride` 是每行之间的字节跨度。 ([Blend2D][3]) + +第二,view 范围必须真的落在 framebuffer 内: + +```text +x >= 0 +y >= 0 + +x + width <= framebuffer.width +y + height <= framebuffer.height +``` + +这个检查我会放在你创建 surface/view 的最外层,而不是每次绘制都检查。 + +第三,也是比较重要的:不要让每个 view 都释放同一个 framebuffer。 + +`create_from_data()` 有: + +```cpp +destroy_func +user_data +``` + +用于外部内存生命周期管理。 ([Blend2D][1]) + +如果你这样干: + +```text +view A destroy → free(framebuffer) +view B destroy → free(framebuffer) +``` + +就是 double-free。 + +你这种结构最简单的是: + +```cpp +image.create_from_data( + ..., + BL_DATA_ACCESS_RW, + nullptr, + nullptr); +``` + +由你的 framebuffer owner 自己负责生命周期: + +```text +Framebuffer + owns memory + + ├── ImageView A + ├── ImageView B + ├── ImageView C + └── ImageView D +``` + +保证: + +```text +所有 ImageView / BLContext 销毁 + ↓ +Framebuffer 才销毁 +``` + +如果必须跨生命周期,就自己搞共享 owner/refcount。 + +第四,不要 weak-copy view 再期待两个 renderer 并行: + +```cpp +BLImage a = make_view(...); +BLImage b = a; +``` + +这不是两个 view。 + +这是共享同一个 Blend2D image data。Blend2D 的 `BLImage` copy 是 weak-copy,会共享底层 image data。 ([Blend2D][1]) + +应该每个 region 独立: + +```cpp +BLImage a; +a.create_from_data(...); + +BLImage b; +b.create_from_data(...); +``` + +第五,我建议你的 view 自己建立局部坐标系。 + +比如 Plot B 在最终 framebuffer: + +```text +x = 960 +y = 0 +w = 960 +h = 540 +``` + +但是它拿到的 BLImage 是: + +```text +0 .. 959 +0 .. 539 +``` + +所以 Plot 内部正常画: + +```cpp +ctx.fill_rect(BLRect(0, 0, 960, 540), ...); +``` + +不需要知道自己在父 framebuffer 里是 `(960, 0)`。 + +这对你的 Scene/Plot 抽象非常干净: + +```text +Parent framebuffer coordinates + ↓ +Surface/View + ↓ +local coordinates + ↓ +Plot renderer +``` + +还有一个容易担心但实际上比较好的点:Blend2D 会裁剪到当前 image target。 + +所以一个 960×540 的 view,即使你的 path 坐标跑到: + +```text +x = 1500 +``` + +Blend2D target 本身只有: + +```text +960 × 540 +``` + +正常 raster 不应该越过 target memory 定义的逻辑宽度去写相邻区域。 + +也就是说 view 本身实际上同时给了你一个天然 clip boundary。 + +所以我会把你的最终结构做成: + +```cpp +class Surface { +public: + BLImage image; + int x; + int y; + int width; + int height; +}; +``` + +然后: + +```text +Framebuffer + │ + ├── Surface Plot0 + │ └── BLImage(external framebuffer region) + │ + ├── Surface Plot1 + │ └── BLImage(external framebuffer region) + │ + └── Surface Plot2 + └── BLImage(external framebuffer region) +``` + +Taskflow: + +```text + framebuffer + + ┌──────────┬──────────┐ + │ Plot A │ Plot B │ + ├──────────┼──────────┤ + │ Plot C │ Plot D │ + └──────────┴──────────┘ + ▲ ▲ + │ │ + Task A Task B +``` + +这种情况下不存在: + +```text +Image A +Image B +Image C + ↓ +composite +``` + +每个任务第一次 rasterize 就已经直接落到最终 framebuffer。 + +不过我要保留一个工程上的判断: + +**这是利用 Blend2D external-image API 实现的物理 framebuffer 分区,并不是 Blend2D 文档明确承诺的“多个 external `BLImage` +指向同一 allocation 的不同区域并发渲染”高级用法。** + +Blend2D 官方承诺的是 external memory image 本身可用,以及一个 context 对其 image data 独占;它不会理解你几个独立 external +image 的 backing allocations 是否重叠。 ([Blend2D][1]) + +所以我会把它当成: + +```text +可实现 ++ +原理成立 ++ +需要你自己严格保证 region 不重叠 ++ +最好做一次 TSAN + 压力测试 + framebuffer guard 测试 +``` + +对于你这种“一个大图天然分成多个 Plot”的场景,我认为值得做,而且比“各画一张临时图最后再 blit”更符合你追求后台性能的方向。 + +[1]: https://blend2d.com/doc/classBLImage.html?utm_source=chatgpt.com "API - Docs - Blend2D" +[2]: https://blend2d.com/doc/classBLContext.html?utm_source=chatgpt.com "API - Docs - Blend2D" +[3]: https://blend2d.com/doc/structBLImageData.html?utm_source=chatgpt.com "API - Docs - Blend2D"