From 65799cfc2a28ecf098c786d9f32819afe8d3743f Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Tue, 25 Aug 2026 21:44:57 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96=E4=BA=86=20=E8=BF=98?= =?UTF-8?q?=E6=98=AF=E5=8D=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kernel/src/kernel/frame_statistics.cpp | 15 +- render_3D/render_3D/Gpu_Completion_State.cpp | 2 +- render_3D/render_3D/Gpu_Completion_State.hpp | 3 +- .../detail/Gpu_Completion_Service.cpp | 747 +++++++++--------- .../detail/Gpu_Completion_Service.hpp | 82 +- web_server/src/Gallery_Video_Stream.cpp | 104 ++- web_server/src/Gallery_WebSocket.cpp | 59 +- web_server/src/Gallery_WebSocket.hpp | 9 +- web_server/src/Graph_WebSocket.cpp | 83 +- web_server/src/Graph_WebSocket.hpp | 10 +- web_server/src/Page_Session.cpp | 80 -- web_server/src/Page_Session.hpp | 33 - web_server/src/Plot.cpp | 393 ++++----- web_server/src/WebRtc_Video_Session.cpp | 176 ++--- web_server/src/Web_Server.cpp | 16 +- web_server/src/detail/Gallery_Frame_Atlas.cpp | 80 +- web_server/src/detail/Gallery_Frame_Clock.cpp | 40 +- web_server/tests/Page_Session_Tests.cpp | 36 - web_server/tests/Sliding_Statistics_Tests.cpp | 15 + webapp_gallery/src/app.tsx | 66 +- webapp_gallery/src/main.tsx | 25 +- 21 files changed, 850 insertions(+), 1224 deletions(-) delete mode 100644 web_server/src/Page_Session.cpp delete mode 100644 web_server/src/Page_Session.hpp delete mode 100644 web_server/tests/Page_Session_Tests.cpp diff --git a/kernel/src/kernel/frame_statistics.cpp b/kernel/src/kernel/frame_statistics.cpp index f6a1188..b637e1d 100644 --- a/kernel/src/kernel/frame_statistics.cpp +++ b/kernel/src/kernel/frame_statistics.cpp @@ -130,8 +130,9 @@ Statistic_State Sliding_Statistics::submit(double value) { if (!std::isfinite(value)) return state_; if (size_ == values_.size()) { if (!trimmed_window_initialized_) { - const double lower = p05_.value(); - const double upper = p95_.value(); + double lower = p05_.value(); + double upper = p95_.value(); + if (lower > upper) std::swap(lower, upper); trimmed_sum_ = 0.0; for (std::size_t index = 0; index < size_; ++index) { trimmed_values_[index] = std::clamp(values_[index], lower, upper); @@ -151,8 +152,11 @@ Statistic_State Sliding_Statistics::submit(double value) { minimum_.submit(value, sequence_, values_.size()); maximum_.submit(value, sequence_, values_.size()); ++sequence_; + double trim_lower = p05_.value(); + double trim_upper = p95_.value(); + if (trim_lower > trim_upper) std::swap(trim_lower, trim_upper); const double trimmed = trimmed_window_initialized_ - ? std::clamp(value, p05_.value(), p95_.value()) : value; + ? std::clamp(value, trim_lower, trim_upper) : value; values_[next_] = value; trimmed_values_[next_] = trimmed; sum_ += value; @@ -161,12 +165,15 @@ Statistic_State Sliding_Statistics::submit(double value) { next_ = (next_ + 1U) % values_.size(); size_ = std::min(size_ + 1U, values_.size()); const double mean = sum_ / static_cast(size_); + const double quantile_50 = p50_.value(); + const double quantile_95 = std::max(quantile_50, p95_.value()); + const double quantile_99 = std::max(quantile_95, p99_.value()); state_ = Statistic_State{ size_, value, minimum_.value(), maximum_.value(), mean, trimmed_sum_ / static_cast(size_), std::sqrt(std::max(0.0, squared_sum_ / static_cast(size_) - mean * mean)), - p50_.value(), p95_.value(), p99_.value()}; + quantile_50, quantile_95, quantile_99}; return state_; } diff --git a/render_3D/render_3D/Gpu_Completion_State.cpp b/render_3D/render_3D/Gpu_Completion_State.cpp index 5a173e0..6091029 100644 --- a/render_3D/render_3D/Gpu_Completion_State.cpp +++ b/render_3D/render_3D/Gpu_Completion_State.cpp @@ -2,7 +2,7 @@ #include "detail/Gpu_Completion_Service.hpp" namespace aethera::render_3d { -Gpu_Completion_State gpu_completion_state() noexcept { +Gpu_Completion_State gpu_completion_state() { return detail::Gpu_Completion_Service::instance().state(); } } diff --git a/render_3D/render_3D/Gpu_Completion_State.hpp b/render_3D/render_3D/Gpu_Completion_State.hpp index 274256c..09eaabf 100644 --- a/render_3D/render_3D/Gpu_Completion_State.hpp +++ b/render_3D/render_3D/Gpu_Completion_State.hpp @@ -27,12 +27,11 @@ struct Gpu_Completion_State : State_Type { std::uint64_t callback_max_ns{}; std::uint64_t callback_failure_count{}; std::uint64_t backpressure_count{}; - std::uint64_t backpressure_wait_ns{}; std::uint64_t fault_count{}; std::uint64_t abandoned_count{}; bool stopping{}; bool operator==(const Gpu_Completion_State&) const = default; }; -[[nodiscard]] Gpu_Completion_State gpu_completion_state() noexcept; +[[nodiscard]] Gpu_Completion_State gpu_completion_state(); } diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp index 7b993b8..e959b63 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp @@ -1,440 +1,439 @@ -#include "Gpu_Completion_Service.hpp" +#include "Gpu_Completion_Service.hpp" #include "Exception.hpp" #include #include #include #include + namespace aethera::render_3d::detail { +namespace { +std::uint64_t elapsed_nanoseconds( + std::chrono::steady_clock::time_point started) noexcept { + const auto elapsed = std::chrono::duration_cast( + std::chrono::steady_clock::now() - started).count(); + return elapsed > 0 ? static_cast(elapsed) : 0ULL; +} + +template +void update_peak(Value& peak, Value value) noexcept { + peak = std::max(peak, value); +} +} + Gpu_Completion_Service& Gpu_Completion_Service::instance() { static Gpu_Completion_Service service; return service; } + Gpu_Completion_Service::Gpu_Completion_Service() { - publish_state(0, 0); - thread_ = std::thread([this] { - run(); - }); + state_.current->capacity = static_cast(default_capacity); + state_.advance(); + thread_ = std::thread([this] { run(); }); } + Gpu_Completion_Service::~Gpu_Completion_Service() { - stopping_.store(true, std::memory_order_release); - wake(); + { + std::lock_guard lock(service_mutex_); + state_.current->stopping = true; + ++wake_generation_; + state_.advance(); + } + wake_condition_.notify_one(); if (thread_.joinable()) thread_.join(); } + Gpu_Completion_Service::Reservation::Reservation( - std::shared_ptr pending) noexcept : pending_(std::move(pending)) {} + std::shared_ptr pending) noexcept + : pending_(std::move(pending)) {} + Gpu_Completion_Service::Reservation::~Reservation() noexcept { - try { cancel(); } + try { + cancel(); + } catch (...) {} } -Gpu_Completion_Service::Reservation::Reservation(Reservation&& other) noexcept : pending_(std::exchange(other.pending_, {})) {} + +Gpu_Completion_Service::Reservation::Reservation( + Reservation&& other) noexcept + : pending_(std::exchange(other.pending_, {})) {} + Gpu_Completion_Service::Prepare_Result::operator bool() const noexcept { return result == Admission_Result::none; } + void Gpu_Completion_Service::Reservation::watch(VkDevice device, VkFence fence) { - if (!pending_ || device == VK_NULL_HANDLE || fence == VK_NULL_HANDLE) throw std::logic_error("GPU completion reservation or fence is invalid"); - auto pending = std::exchange(pending_, {}); - auto* const service = pending->service; - { - std::lock_guard lock(pending->mutex); - if (pending->status != Pending_Fence::Status::reserved) throw std::logic_error("GPU completion reservation is not reserved"); - pending->device = device; - pending->fence = fence; - // This timestamp is part of correctness, not only observability: it - // bounds the lifetime of a submission whose fence never signals. - pending->watched_at = std::chrono::steady_clock::now(); - const std::size_t watched = - service->watched_.fetch_add(1, std::memory_order_relaxed) + 1; - update_peak(service->peak_watched_, watched); - pending->status = Pending_Fence::Status::watched; - } - service->wake(); + if (!pending_ || device == VK_NULL_HANDLE || fence == VK_NULL_HANDLE) + throw std::logic_error( + "GPU completion reservation or fence is invalid"); + pending_->service->watch(pending_, device, fence); + pending_.reset(); } + void Gpu_Completion_Service::Reservation::cancel() { if (!pending_) return; - auto pending = std::exchange(pending_, {}); - cancel_reserved(pending); - pending->service->wake(); + pending_->service->cancel(pending_); + pending_.reset(); } -void Gpu_Completion_Service::update_peak(std::atomic_size_t& peak, - std::size_t value) noexcept { - std::size_t current = peak.load(std::memory_order_relaxed); - while (current < value && - !peak.compare_exchange_weak(current, value, - std::memory_order_relaxed)) {} -} -void Gpu_Completion_Service::update_max(std::atomic_uint64_t& maximum, - std::uint64_t value) noexcept { - std::uint64_t current = maximum.load(std::memory_order_relaxed); - while (current < value && - !maximum.compare_exchange_weak(current, value, - std::memory_order_relaxed)) {} -} -void Gpu_Completion_Service::cancel_reserved( - const std::shared_ptr& pending) { - if (!pending) return; - std::lock_guard lock(pending->mutex); - if (pending->status == Pending_Fence::Status::reserved) pending->status = Pending_Fence::Status::canceled; -} -bool Gpu_Completion_Service::acquire_slot() noexcept { - std::size_t current = in_flight_.load(std::memory_order_relaxed); - while (current < static_cast(default_capacity)) { - if (in_flight_.compare_exchange_weak( - current, current + 1, - std::memory_order_acq_rel, - std::memory_order_relaxed)) { - update_peak(peak_in_flight_, current + 1); - return true; - } + +void Gpu_Completion_Service::watch( + const std::shared_ptr& pending, + VkDevice device, VkFence fence) { + { + std::lock_guard lock(service_mutex_); + if (!pending || pending->service != this || + pending->status != Pending_Fence::Status::reserved) + throw std::logic_error( + "GPU completion reservation is not reserved"); + pending->device = device; + pending->fence = fence; + pending->watched_at = std::chrono::steady_clock::now(); + pending->status = Pending_Fence::Status::watched; + ++state_.current->watched; + update_peak(state_.current->peak_watched, + state_.current->watched); + ++wake_generation_; } - backpressure_count_.fetch_add(1, std::memory_order_relaxed); - return false; + wake_condition_.notify_one(); } -void Gpu_Completion_Service::release_slot() noexcept { - in_flight_.fetch_sub(1, std::memory_order_relaxed); + +void Gpu_Completion_Service::cancel( + const std::shared_ptr& pending) { + { + std::lock_guard lock(service_mutex_); + if (!pending || pending->service != this) return; + if (pending->status == Pending_Fence::Status::reserved) + pending->status = Pending_Fence::Status::canceled; + ++wake_generation_; + } + wake_condition_.notify_one(); } + Gpu_Completion_Service::Prepare_Result Gpu_Completion_Service::prepare( Completion completion, Exception_Handler on_exception, bool observe) { - if (!completion) throw std::invalid_argument("GPU completion callback is empty"); - if (!on_exception) throw std::invalid_argument("GPU completion exception handler is empty"); - if (stopping_.load(std::memory_order_acquire)) return {{}, Admission_Result::stopping}; + if (!completion) + throw std::invalid_argument("GPU completion callback is empty"); + if (!on_exception) + throw std::invalid_argument( + "GPU completion exception handler is empty"); + auto pending = std::make_shared(); pending->completion = std::move(completion); pending->on_exception = std::move(on_exception); pending->observe = observe; pending->service = this; - /* 唯一 GPU Submit 域绝不等待 completion 容量。满载时把背压作为 - * 准入结果立即反馈给 Scene,由上游下一帧策略自然重试。 */ - if (!acquire_slot()) - return {{}, Admission_Result::capacity_exhausted}; - if (stopping_.load(std::memory_order_acquire)) { - release_slot(); - return {{}, Admission_Result::stopping}; - } - try { - std::lock_guard lock(pending_mutex_); + { + std::lock_guard lock(service_mutex_); + auto& state = *state_.current; + if (state.stopping) + return {{}, Admission_Result::stopping}; + if (state.in_flight >= state.capacity) { + ++state.backpressure_count; + return {{}, Admission_Result::capacity_exhausted}; + } pending_.push_back(pending); + ++state.in_flight; + update_peak(state.peak_in_flight, state.in_flight); + ++state.reservation_count; + ++wake_generation_; } - catch (...) { - release_slot(); - throw; - } - reservation_count_.fetch_add(1, std::memory_order_relaxed); - wake(); + wake_condition_.notify_one(); return {Reservation(std::move(pending)), Admission_Result::none}; } -Gpu_Completion_State Gpu_Completion_Service::state() const noexcept { - std::lock_guard lock(state_mutex_); + +Gpu_Completion_State Gpu_Completion_Service::state() const { + std::lock_guard lock(service_mutex_); return *state_.pending; } -void Gpu_Completion_Service::publish_state(std::size_t active_fences, - std::size_t pending_fences) noexcept { - std::lock_guard lock(state_mutex_); - auto& state = *state_.current; - state.capacity = static_cast(default_capacity); - state.in_flight = in_flight_.load(std::memory_order_relaxed); - state.peak_in_flight = peak_in_flight_.load(std::memory_order_relaxed); - state.watched = watched_.load(std::memory_order_relaxed); - state.peak_watched = peak_watched_.load(std::memory_order_relaxed); - state.active_fences = active_fences; - state.pending_fences = pending_fences; - state.reservation_count = reservation_count_.load(std::memory_order_relaxed); - state.completion_count = completion_count_.load(std::memory_order_relaxed); - state.cancellation_count = cancellation_count_.load(std::memory_order_relaxed); - state.fence_probe_count = fence_probe_count_.load(std::memory_order_relaxed); - state.fence_wait_count = fence_wait_count_.load(std::memory_order_relaxed); - state.fence_wait_timeout_count = fence_wait_timeout_count_.load(std::memory_order_relaxed); - state.fence_wait_total_ns = fence_wait_total_ns_.load(std::memory_order_relaxed); - state.fence_wait_max_ns = fence_wait_max_ns_.load(std::memory_order_relaxed); - state.callback_total_ns = callback_total_ns_.load(std::memory_order_relaxed); - state.callback_max_ns = callback_max_ns_.load(std::memory_order_relaxed); - state.callback_failure_count = callback_failure_count_.load(std::memory_order_relaxed); - state.backpressure_count = backpressure_count_.load(std::memory_order_relaxed); - state.backpressure_wait_ns = backpressure_wait_ns_.load(std::memory_order_relaxed); - state.fault_count = fault_count_.load(std::memory_order_relaxed); - state.abandoned_count = abandoned_count_.load(std::memory_order_relaxed); - state.stopping = stopping_.load(std::memory_order_relaxed); + +void Gpu_Completion_Service::publish_state( + std::size_t active_fences, std::size_t pending_fences) { + std::lock_guard lock(service_mutex_); + state_.current->active_fences = active_fences; + state_.current->pending_fences = pending_fences; state_.advance(); } -void Gpu_Completion_Service::wake() noexcept { - wake_generation_.fetch_add(1, std::memory_order_release); - wake_condition_.notify_one(); -} + void Gpu_Completion_Service::run() noexcept { struct Device_Fences { VkDevice device{VK_NULL_HANDLE}; - std::vector fences; + std::vector fences{}; }; std::vector> active; + try { - active.reserve(static_cast(default_capacity)); - std::size_t wait_group_index{}; - auto last_state_publication = std::chrono::steady_clock::time_point{}; - const auto wait_age_ns = [](std::chrono::steady_clock::time_point started) { - const auto elapsed = std::chrono::duration_cast( - std::chrono::steady_clock::now() - started) - .count(); - return elapsed > 0 ? static_cast(elapsed) : 0ULL; - }; - const auto finish = [this, &wait_age_ns]( - const std::shared_ptr& pending, - VkResult result, Completion_Error error) noexcept { - struct Slot_Release { - Gpu_Completion_Service* service; - ~Slot_Release() { service->release_slot(); } - } slot_release{this}; - watched_.fetch_sub(1, std::memory_order_relaxed); - try { - Completion completion; - Exception_Handler on_exception; - std::chrono::steady_clock::time_point watched_at{}; - bool observe{}; - { - std::lock_guard lock(pending->mutex); - completion = std::move(pending->completion); - on_exception = std::move(pending->on_exception); - watched_at = pending->watched_at; - observe = pending->observe; - pending->status = Pending_Fence::Status::canceled; - } - Result completion_result; - completion_result.error = error; - completion_result.vulkan_result = result; - if (observe) completion_result.wait_duration_ns = wait_age_ns(watched_at); - if (error != Completion_Error::none) { - fault_count_.fetch_add(1, std::memory_order_relaxed); - abandoned_count_.fetch_add(1, std::memory_order_relaxed); - } - const auto callback_started = std::chrono::steady_clock::now(); - try { completion(std::move(completion_result)); } - catch (...) { - callback_failure_count_.fetch_add(1, std::memory_order_relaxed); - try { - on_exception(contextual_exception( - "delivering GPU completion", std::current_exception())); - } - catch (...) {} - } - const auto callback_elapsed = std::chrono::duration_cast( - std::chrono::steady_clock::now() - callback_started).count(); - if (callback_elapsed > 0) { - const auto elapsed = static_cast(callback_elapsed); - callback_total_ns_.fetch_add(elapsed, std::memory_order_relaxed); - update_max(callback_max_ns_, elapsed); - } - completion_count_.fetch_add(1, std::memory_order_relaxed); - } - catch (...) { - try { - std::lock_guard lock(pending->mutex); - if (pending->on_exception) - pending->on_exception(std::current_exception()); - pending->status = Pending_Fence::Status::canceled; - } - catch (...) {} - } - }; - for (;;) { - const std::uint64_t wake_generation = - wake_generation_.load(std::memory_order_acquire); - { std::lock_guard lock(pending_mutex_); while (!pending_.empty()) { active.push_back(std::move(pending_.front())); pending_.pop_front(); } } - std::vector groups; - std::size_t reserved_count{}; - for (auto iterator = active.begin(); iterator != active.end();) { - Pending_Fence::Status status; - VkDevice device{VK_NULL_HANDLE}; - VkFence fence{VK_NULL_HANDLE}; - { - std::lock_guard lock((*iterator)->mutex); - status = (*iterator)->status; - device = (*iterator)->device; - fence = (*iterator)->fence; - } - if (status == Pending_Fence::Status::canceled || - (status == Pending_Fence::Status::reserved && - stopping_.load(std::memory_order_acquire))) { - cancel_reserved(*iterator); - iterator = active.erase(iterator); - cancellation_count_.fetch_add(1, std::memory_order_relaxed); - release_slot(); - continue; - } - if (status == Pending_Fence::Status::reserved) ++reserved_count; - if (status == Pending_Fence::Status::watched) { - auto group = std::find_if( - groups.begin(), groups.end(), - [device](const Device_Fences& item) { - return item.device == device; - }); - if (group == groups.end()) { - groups.push_back(Device_Fences{device, {}}); - group = groups.end() - 1; - } - group->fences.push_back(fence); - } - ++iterator; - } - const auto publication_time = std::chrono::steady_clock::now(); - if (last_state_publication == std::chrono::steady_clock::time_point{} || - publication_time - last_state_publication >= state_publication_interval) { - // 域内只按诊断粒度发布双缓冲 State;逐 fence 热路径只更新原子计数。 - publish_state(active.size(), reserved_count); - last_state_publication = publication_time; - } - bool pending_empty; - { std::lock_guard lock(pending_mutex_); pending_empty = pending_.empty(); } - if (stopping_.load(std::memory_order_acquire) && active.empty() && pending_empty) { - publish_state(0, 0); - return; - } - // Probe every watched fence first. A permanently unsignaled fence is - // converted into a logical failure after a bounded interval. The - // submission is explicitly marked abandoned so its owner can - // quarantine, rather than recycle, the referenced GPU resources. - bool completed_any = false; - for (auto iterator = active.begin(); iterator != active.end();) { - VkDevice device{VK_NULL_HANDLE}; - VkFence fence{VK_NULL_HANDLE}; - Pending_Fence::Status status; - std::chrono::steady_clock::time_point watched_at{}; - { - std::lock_guard lock((*iterator)->mutex); - status = (*iterator)->status; - device = (*iterator)->device; - fence = (*iterator)->fence; - watched_at = (*iterator)->watched_at; - } - if (status != Pending_Fence::Status::watched) { - ++iterator; - continue; - } - if (stopping_.load(std::memory_order_acquire)) { - auto pending = *iterator; - iterator = active.erase(iterator); - finish(pending, VK_TIMEOUT, Completion_Error::fence_abandoned); - completed_any = true; - continue; - } - if (wait_age_ns(watched_at) >= maximum_fence_age_ns) { - auto pending = *iterator; - iterator = active.erase(iterator); - finish(pending, VK_TIMEOUT, Completion_Error::fence_abandoned); - completed_any = true; - continue; - } - fence_probe_count_.fetch_add(1, std::memory_order_relaxed); - const VkResult result = vkGetFenceStatus(device, fence); - if (result == VK_NOT_READY) { - ++iterator; - continue; - } - auto pending = *iterator; - iterator = active.erase(iterator); - finish(pending, result, result == VK_SUCCESS - ? Completion_Error::none - : Completion_Error::vulkan_failure); - completed_any = true; - } - if (completed_any) continue; - if (groups.empty()) { - std::unique_lock lock(wait_mutex_); - if (wake_generation_.load(std::memory_order_acquire) == - wake_generation) { - wake_condition_.wait(lock, [this, wake_generation] { - return wake_generation_.load(std::memory_order_acquire) != - wake_generation; - }); - } - continue; - } - wait_group_index %= groups.size(); - const Device_Fences& group = groups[wait_group_index++]; - const auto wait_started = std::chrono::steady_clock::now(); - fence_wait_count_.fetch_add(1, std::memory_order_relaxed); - const VkResult wait_result = vkWaitForFences( - group.device, static_cast(group.fences.size()), - group.fences.data(), VK_FALSE, fence_wait_timeout_ns); - const auto wait_elapsed = std::chrono::duration_cast( - std::chrono::steady_clock::now() - wait_started).count(); - if (wait_elapsed > 0) { - const auto elapsed = static_cast(wait_elapsed); - fence_wait_total_ns_.fetch_add(elapsed, std::memory_order_relaxed); - update_max(fence_wait_max_ns_, elapsed); - } - if (wait_result == VK_TIMEOUT) - fence_wait_timeout_count_.fetch_add(1, std::memory_order_relaxed); - for (auto iterator = active.begin(); iterator != active.end();) { - VkDevice device{VK_NULL_HANDLE}; - VkFence fence{VK_NULL_HANDLE}; - Pending_Fence::Status status; - { - std::lock_guard lock((*iterator)->mutex); - status = (*iterator)->status; - device = (*iterator)->device; - fence = (*iterator)->fence; - } - if (status != Pending_Fence::Status::watched || - device != group.device) { - ++iterator; - continue; - } - VkResult result = wait_result; - if (wait_result == VK_SUCCESS || wait_result == VK_TIMEOUT) { - fence_probe_count_.fetch_add(1, std::memory_order_relaxed); - result = vkGetFenceStatus(device, fence); - } - if (result == VK_NOT_READY) { - ++iterator; - continue; - } - auto pending = *iterator; - iterator = active.erase(iterator); - finish(pending, result, result == VK_SUCCESS - ? Completion_Error::none - : Completion_Error::vulkan_failure); - } - } - } - catch (...) { - stopping_.store(true, std::memory_order_release); - const auto service_failure = std::current_exception(); - const auto fail_pending = [this, &service_failure]( - const std::shared_ptr& pending) noexcept { - if (!pending) return; + active.reserve(static_cast(default_capacity)); + std::size_t wait_group_index{}; + auto last_state_publication = + std::chrono::steady_clock::time_point{}; + + const auto finish = [this]( + const std::shared_ptr& pending, + VkResult vulkan_result, Completion_Error error) noexcept { + Completion completion; Exception_Handler on_exception; - bool watched{}; - try { - std::lock_guard lock(pending->mutex); - watched = pending->status == Pending_Fence::Status::watched; + Result result{}; + { + std::lock_guard lock(service_mutex_); + if (pending->status != Pending_Fence::Status::watched) + return; + completion = std::move(pending->completion); on_exception = std::move(pending->on_exception); - pending->completion = {}; + result.error = error; + result.vulkan_result = vulkan_result; + if (pending->observe) + result.wait_duration_ns = + elapsed_nanoseconds(pending->watched_at); pending->status = Pending_Fence::Status::canceled; + auto& state = *state_.current; + if (state.watched != 0) --state.watched; + if (error != Completion_Error::none) { + ++state.fault_count; + if (error == Completion_Error::fence_abandoned) + ++state.abandoned_count; + } } - catch (...) {} - if (watched) - watched_.fetch_sub(1, std::memory_order_relaxed); - cancellation_count_.fetch_add(1, std::memory_order_relaxed); - if (on_exception) { - try { on_exception(service_failure); } + + const auto callback_started = + std::chrono::steady_clock::now(); + std::exception_ptr callback_failure; + try { + completion(std::move(result)); + } + catch (...) { + callback_failure = std::current_exception(); + } + const auto callback_ns = + elapsed_nanoseconds(callback_started); + { + std::lock_guard lock(service_mutex_); + auto& state = *state_.current; + state.callback_total_ns += callback_ns; + update_peak(state.callback_max_ns, callback_ns); + if (callback_failure) ++state.callback_failure_count; + ++state.completion_count; + if (state.in_flight != 0) --state.in_flight; + } + if (callback_failure) { + try { + on_exception(contextual_exception( + "delivering GPU completion", + std::move(callback_failure))); + } catch (...) {} } - release_slot(); }; - for (const auto& pending : active) fail_pending(pending); + for (;;) { - std::shared_ptr pending; - try { - std::lock_guard lock(pending_mutex_); - if (pending_.empty()) break; - pending = std::move(pending_.front()); + std::uint64_t wake_generation{}; + bool stopping{}; + std::size_t reserved_count{}; + std::vector groups; + { + std::lock_guard lock(service_mutex_); + wake_generation = wake_generation_; + stopping = state_.current->stopping; + while (!pending_.empty()) { + active.push_back(std::move(pending_.front())); + pending_.pop_front(); + } + for (auto iterator = active.begin(); + iterator != active.end();) { + const auto status = (*iterator)->status; + if (status == Pending_Fence::Status::canceled || + (status == Pending_Fence::Status::reserved && + stopping)) { + iterator = active.erase(iterator); + ++state_.current->cancellation_count; + if (state_.current->in_flight != 0) + --state_.current->in_flight; + continue; + } + if (status == Pending_Fence::Status::reserved) { + ++reserved_count; + ++iterator; + continue; + } + auto group = std::ranges::find( + groups, (*iterator)->device, + &Device_Fences::device); + if (group == groups.end()) { + groups.push_back( + Device_Fences{(*iterator)->device, {}}); + group = groups.end() - 1; + } + group->fences.push_back((*iterator)->fence); + ++iterator; + } + } + + const auto publication_time = + std::chrono::steady_clock::now(); + if (last_state_publication == + std::chrono::steady_clock::time_point{} || + publication_time - last_state_publication >= + state_publication_interval) { + publish_state(active.size(), reserved_count); + last_state_publication = publication_time; + } + + bool pending_empty{}; + { + std::lock_guard lock(service_mutex_); + stopping = state_.current->stopping; + pending_empty = pending_.empty(); + } + if (stopping && active.empty() && pending_empty) { + publish_state(0, 0); + return; + } + + bool completed_any{}; + for (auto iterator = active.begin(); + iterator != active.end();) { + VkDevice device{VK_NULL_HANDLE}; + VkFence fence{VK_NULL_HANDLE}; + std::chrono::steady_clock::time_point watched_at{}; + Pending_Fence::Status status{}; + { + std::lock_guard lock(service_mutex_); + status = (*iterator)->status; + device = (*iterator)->device; + fence = (*iterator)->fence; + watched_at = (*iterator)->watched_at; + stopping = state_.current->stopping; + } + if (status != Pending_Fence::Status::watched) { + ++iterator; + continue; + } + if (stopping || + elapsed_nanoseconds(watched_at) >= + maximum_fence_age_ns) { + auto pending = *iterator; + iterator = active.erase(iterator); + finish(pending, VK_TIMEOUT, + Completion_Error::fence_abandoned); + completed_any = true; + continue; + } + { + std::lock_guard lock(service_mutex_); + ++state_.current->fence_probe_count; + } + const VkResult result = vkGetFenceStatus(device, fence); + if (result == VK_NOT_READY) { + ++iterator; + continue; + } + auto pending = *iterator; + iterator = active.erase(iterator); + finish(pending, result, + result == VK_SUCCESS + ? Completion_Error::none + : Completion_Error::vulkan_failure); + completed_any = true; + } + if (completed_any) continue; + + if (groups.empty()) { + std::unique_lock lock(service_mutex_); + wake_condition_.wait(lock, [this, wake_generation] { + return wake_generation_ != wake_generation || + state_.current->stopping; + }); + continue; + } + + wait_group_index %= groups.size(); + const auto& group = groups[wait_group_index++]; + const auto wait_started = std::chrono::steady_clock::now(); + { + std::lock_guard lock(service_mutex_); + ++state_.current->fence_wait_count; + } + const VkResult wait_result = vkWaitForFences( + group.device, + static_cast(group.fences.size()), + group.fences.data(), VK_FALSE, + fence_wait_timeout_ns); + const auto wait_ns = elapsed_nanoseconds(wait_started); + { + std::lock_guard lock(service_mutex_); + auto& state = *state_.current; + state.fence_wait_total_ns += wait_ns; + update_peak(state.fence_wait_max_ns, wait_ns); + if (wait_result == VK_TIMEOUT) + ++state.fence_wait_timeout_count; + } + + for (auto iterator = active.begin(); + iterator != active.end();) { + VkDevice device{VK_NULL_HANDLE}; + VkFence fence{VK_NULL_HANDLE}; + Pending_Fence::Status status{}; + { + std::lock_guard lock(service_mutex_); + status = (*iterator)->status; + device = (*iterator)->device; + fence = (*iterator)->fence; + } + if (status != Pending_Fence::Status::watched || + device != group.device) { + ++iterator; + continue; + } + VkResult result = wait_result; + if (wait_result == VK_SUCCESS || + wait_result == VK_TIMEOUT) { + { + std::lock_guard lock(service_mutex_); + ++state_.current->fence_probe_count; + } + result = vkGetFenceStatus(device, fence); + } + if (result == VK_NOT_READY) { + ++iterator; + continue; + } + auto pending = *iterator; + iterator = active.erase(iterator); + finish(pending, result, + result == VK_SUCCESS + ? Completion_Error::none + : Completion_Error::vulkan_failure); + } + } + } + catch (...) { + const auto service_failure = std::current_exception(); + std::vector> failed = + std::move(active); + { + std::lock_guard lock(service_mutex_); + while (!pending_.empty()) { + failed.push_back(std::move(pending_.front())); pending_.pop_front(); } - catch (...) { break; } - fail_pending(pending); + auto& state = *state_.current; + state.stopping = true; + for (const auto& pending : failed) { + if (pending->status == Pending_Fence::Status::watched && + state.watched != 0) + --state.watched; + pending->status = Pending_Fence::Status::canceled; + ++state.cancellation_count; + if (state.in_flight != 0) --state.in_flight; + } + state.active_fences = 0; + state.pending_fences = 0; + state_.advance(); + } + for (auto& pending : failed) { + if (!pending->on_exception) continue; + try { pending->on_exception(service_failure); } + catch (...) {} } - publish_state(0, 0); } } } - - diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp index 1f06e05..8ab5105 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp @@ -2,7 +2,6 @@ #include "../Gpu_Completion_State.hpp" #include #include -#include #include #include #include @@ -28,9 +27,9 @@ public: vulkan_failure }; struct Result { - Completion_Error error{}; /* 归一化完成结果。 */ - VkResult vulkan_result{VK_SUCCESS}; /* Vulkan 原始结果码。 */ - std::uint64_t wait_duration_ns{}; /* fence 等待时间,单位为纳秒。 */ + Completion_Error error{}; /* 归一化完成结果。 */ + VkResult vulkan_result{VK_SUCCESS}; /* Vulkan 原始结果码。 */ + std::uint64_t wait_duration_ns{}; /* fence 等待时间,单位为纳秒。 */ }; using Completion = std::function; using Exception_Handler = std::function; @@ -46,12 +45,12 @@ public: private: explicit Reservation(std::shared_ptr pending) noexcept; void cancel(); - std::shared_ptr pending_; /* 尚未 watch 或 cancel 的准入记录。 */ + std::shared_ptr pending_; /* 尚未 watch 或 cancel 的准入记录。 */ friend class Gpu_Completion_Service; }; struct Prepare_Result { - Reservation reservation; /* 成功时返回的 fence reservation。 */ - Admission_Result result{}; /* 准入结果。 */ + Reservation reservation; /* 成功时返回的 fence reservation。 */ + Admission_Result result{}; /* 准入结果。 */ [[nodiscard]] explicit operator bool() const noexcept; }; static Gpu_Completion_Service& instance(); @@ -60,7 +59,7 @@ public: [[nodiscard]] Prepare_Result prepare(Completion completion, Exception_Handler on_exception, bool observe); - [[nodiscard]] Gpu_Completion_State state() const noexcept; + [[nodiscard]] Gpu_Completion_State state() const; private: struct Pending_Fence { enum class Status { @@ -68,61 +67,34 @@ private: watched, canceled }; - std::mutex mutex; /* 保护本条记录的状态和回调移动。 */ - VkDevice device{VK_NULL_HANDLE}; /* fence 所属 Vulkan Device。 */ - VkFence fence{VK_NULL_HANDLE}; /* 受监视的 Vulkan fence。 */ - Completion completion; /* 完成或隔离后的交付回调。 */ - Exception_Handler on_exception; /* 回调异常的隔离入口。 */ - std::chrono::steady_clock::time_point watched_at{}; /* 开始监视的单调时钟时刻。 */ - Gpu_Completion_Service* service{}; /* 不拥有的服务实例。 */ - Status status{Status::reserved}; /* reservation 生命周期状态。 */ - bool observe{}; /* 是否采集 fence 等待耗时。 */ + VkDevice device{VK_NULL_HANDLE}; /* fence 所属 Vulkan Device。 */ + VkFence fence{VK_NULL_HANDLE}; /* 受监视的 Vulkan fence。 */ + Completion completion; /* 完成或隔离后的交付回调。 */ + Exception_Handler on_exception; /* 回调异常的隔离入口。 */ + std::chrono::steady_clock::time_point watched_at{}; /* 开始监视的单调时钟时刻。 */ + Gpu_Completion_Service* service{}; /* 不拥有的服务实例。 */ + Status status{Status::reserved}; /* reservation 生命周期状态。 */ + bool observe{}; /* 是否采集 fence 等待耗时。 */ }; Gpu_Completion_Service(); ~Gpu_Completion_Service(); - static void update_peak(std::atomic_size_t& peak, std::size_t value) noexcept; - static void update_max(std::atomic_uint64_t& maximum, - std::uint64_t value) noexcept; - static void cancel_reserved(const std::shared_ptr& pending); - [[nodiscard]] bool acquire_slot() noexcept; - void release_slot() noexcept; + void watch(const std::shared_ptr& pending, + VkDevice device, VkFence fence); + void cancel(const std::shared_ptr& pending); void publish_state(std::size_t active_fences, - std::size_t pending_fences) noexcept; - void wake() noexcept; + std::size_t pending_fences); void run() noexcept; static constexpr std::ptrdiff_t default_capacity = 1024; static constexpr std::uint64_t fence_wait_timeout_ns = 1'000'000; static constexpr std::uint64_t maximum_fence_age_ns = 30'000'000'000ULL; static constexpr auto state_publication_interval = std::chrono::milliseconds(100); - std::mutex pending_mutex_; /* 保护新登记 fence 队列。 */ - std::deque> pending_; /* 完成线程尚未分组的 fence。 */ - std::mutex wait_mutex_; /* 保护完成线程的条件等待。 */ - std::condition_variable wake_condition_; /* 新 fence 或停止请求的唤醒源。 */ - std::atomic_uint64_t wake_generation_{}; /* 防止丢失唤醒的版本。 */ - std::atomic_size_t in_flight_{}; /* 当前 reservation 数。 */ - std::atomic_size_t peak_in_flight_{}; /* 历史最大 reservation 数。 */ - std::atomic_size_t watched_{}; /* 当前受监视 fence 数。 */ - std::atomic_size_t peak_watched_{}; /* 历史最大受监视 fence 数。 */ - std::atomic_uint64_t reservation_count_{}; - std::atomic_uint64_t completion_count_{}; - std::atomic_uint64_t cancellation_count_{}; - std::atomic_uint64_t fence_probe_count_{}; - std::atomic_uint64_t fence_wait_count_{}; - std::atomic_uint64_t fence_wait_timeout_count_{}; - std::atomic_uint64_t fence_wait_total_ns_{}; - std::atomic_uint64_t fence_wait_max_ns_{}; - std::atomic_uint64_t callback_total_ns_{}; - std::atomic_uint64_t callback_max_ns_{}; - std::atomic_uint64_t callback_failure_count_{}; - std::atomic_uint64_t backpressure_count_{}; /* 准入背压累计次数。 */ - std::atomic_uint64_t backpressure_wait_ns_{}; /* 准入背压累计纳秒。 */ - std::atomic_uint64_t fault_count_{}; /* 完成故障累计次数。 */ - std::atomic_uint64_t abandoned_count_{}; /* 放弃 fence 累计次数。 */ - std::atomic_bool stopping_{}; /* 服务是否正在停止。 */ - mutable std::mutex state_mutex_; /* 只保护域级 State 的指针交换。 */ - double_buffer::Publish_Double_Buffer state_{}; - std::thread thread_; /* 专用 fence 完成线程。 */ + /* GPU completion 的队列、reservation 生命周期和 State 指针交换是 + * 同一个一致性域;这把锁是该域唯一的外部同步边界。 */ + mutable std::mutex service_mutex_{}; + std::condition_variable wake_condition_{}; /* 专用完成线程的休眠唤醒源。 */ + std::uint64_t wake_generation_{}; /* service_mutex_ 下的防丢唤醒版本。 */ + std::deque> pending_{}; /* 尚未转入完成线程本地集合的记录。 */ + double_buffer::Publish_Double_Buffer state_{}; /* GPU completion 运行状态的唯一权威源。 */ + std::thread thread_; /* 专用 fence 完成线程。 */ }; } - - diff --git a/web_server/src/Gallery_Video_Stream.cpp b/web_server/src/Gallery_Video_Stream.cpp index 581d0e7..3043723 100644 --- a/web_server/src/Gallery_Video_Stream.cpp +++ b/web_server/src/Gallery_Video_Stream.cpp @@ -10,8 +10,6 @@ #include #include #include -#include -#include #include #include #include @@ -123,8 +121,7 @@ struct Gallery_Video_Stream::Private { Sliding_Statistics compose_ms{600}; Sliding_Statistics encode_ms{600}; Sliding_Statistics publish_ms{600}; - mutable std::shared_mutex diagnostics_exchange_mutex; - double_buffer::Double_Buffer diagnostics{}; + std::atomic> diagnostics{}; } state{}; std::vector sources{}; /* 已按业务标识排序的稳定图集来源。 */ @@ -141,9 +138,8 @@ struct Gallery_Video_Stream::Private { std::shared_ptr video_readiness{}; /* 编码前查询 WebRTC 三槽发送窗口。 */ std::shared_ptr video_diagnostics{}; }; - mutable std::mutex consumers_mutex; - Stream_Id next_consumer_id{1}; /* 在互斥量内区分重连前后的独占消费者。 */ - std::optional consumer{}; /* 独占网页对应的唯一媒体消费者。 */ + std::atomic_uint64_t next_consumer_id{1}; + std::atomic> consumer{}; /* 最新连接原子接管旧连接。 */ std::atomic_bool stopping{}; /* 拒绝关闭开始后的新 tick、帧与订阅。 */ std::atomic_bool key_frame_requested{}; /* 下一次实际编码前消费的关键帧请求。 */ std::atomic_uint64_t skipped_encode_ticks{}; /* 被最新帧策略合并的页面编码 tick 总数。 */ @@ -177,16 +173,13 @@ struct Gallery_Video_Stream::Private { [[nodiscard]] std::shared_ptr consumer_handler() const { - std::lock_guard lock(consumers_mutex); - return consumer ? consumer->handler : nullptr; + const auto current = consumer.load(std::memory_order_acquire); + return current ? current->handler : nullptr; } [[nodiscard]] bool consumer_accepts_video() const noexcept { - std::shared_ptr readiness; - { - std::lock_guard lock(consumers_mutex); - if (consumer) readiness = consumer->video_readiness; - } + const auto current = consumer.load(std::memory_order_acquire); + const auto readiness = current ? current->video_readiness : nullptr; if (!readiness) return false; try { return (*readiness)(); } catch (...) { return false; } @@ -204,11 +197,11 @@ struct Gallery_Video_Stream::Private { return; } catch (...) {} - { - std::lock_guard lock(consumers_mutex); - if (consumer && consumer->handler == handler) - consumer.reset(); - } + auto current = consumer.load(std::memory_order_acquire); + if (current && current->handler == handler) + static_cast(consumer.compare_exchange_strong( + current, {}, std::memory_order_acq_rel, + std::memory_order_acquire)); if (clock) clock->stop(); } catch (...) {} @@ -345,9 +338,9 @@ struct Gallery_Video_Stream::Private { metric_started = now; metric_encoded_start = encoded_frame_count; metric_clock_start = clock_total; - *state.diagnostics.current = std::move(output); - std::unique_lock lock(state.diagnostics_exchange_mutex); - state.diagnostics.advance(); + state.diagnostics.store( + std::make_shared(std::move(output)), + std::memory_order_release); } void compose_media_frame() { @@ -502,25 +495,27 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( if (!video_diagnostics) throw std::invalid_argument( "gallery video subscription requires WebRTC diagnostics"); - Stream_Id id{}; - { - std::lock_guard lock(d->consumers_mutex); - id = d->next_consumer_id++; - // 每个媒体组只保留最新浏览器连接;新页面原子替换旧回调,旧连接随后 - // unsubscribe(id) 时不会误删新连接。连接接管不依赖页面租约或关闭时序。 - d->consumer = Private::Consumer{ + const auto id = d->next_consumer_id.fetch_add( + 1, std::memory_order_relaxed); + // 每个媒体组只保留最新浏览器连接;新页面原子替换旧回调,旧连接随后 + // unsubscribe(id) 时不会误删新连接。连接接管不依赖页面租约或关闭时序。 + auto consumer = std::make_shared( + Private::Consumer{ id, std::make_shared(std::move(handler)), std::make_shared( std::move(video_readiness)), std::make_shared( - std::move(video_diagnostics))}; - } + std::move(video_diagnostics))}); + d->consumer.store(consumer, std::memory_order_release); try { d->clock->start(); } catch (...) { - std::lock_guard lock(d->consumers_mutex); - if (d->consumer && d->consumer->id == id) d->consumer.reset(); + auto current = d->consumer.load(std::memory_order_acquire); + if (current && current->id == id) + static_cast(d->consumer.compare_exchange_strong( + current, {}, std::memory_order_acq_rel, + std::memory_order_acquire)); throw; } return id; @@ -528,11 +523,13 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( void Gallery_Video_Stream::unsubscribe(Stream_Id stream) { bool removed{}; - { - std::lock_guard lock(d->consumers_mutex); - if (d->consumer && d->consumer->id == stream) { - d->consumer.reset(); + auto current = d->consumer.load(std::memory_order_acquire); + while (current && current->id == stream) { + if (d->consumer.compare_exchange_weak( + current, {}, std::memory_order_acq_rel, + std::memory_order_acquire)) { removed = true; + break; } } if (removed) d->clock->stop(); @@ -553,11 +550,7 @@ void Gallery_Video_Stream::shutdown() noexcept { catch (...) {} source.stream = 0; } - try { - std::lock_guard lock(d->consumers_mutex); - d->consumer.reset(); - } - catch (...) {} + d->consumer.store({}, std::memory_order_release); } std::string Gallery_Video_Stream::layout_description() const { @@ -583,22 +576,17 @@ std::string Gallery_Video_Stream::layout_description() const { } nlohmann::json Gallery_Video_Stream::diagnostics() const { - nlohmann::json output; - { - std::shared_lock lock(d->state.diagnostics_exchange_mutex); - output = !d->state.diagnostics.pending->is_null() - ? *d->state.diagnostics.pending - : nlohmann::json{{"kind", "gallery_metrics"}, - {"protocol", "aethera.gallery.video"}, - {"version", 2}, - {"sources", nlohmann::json::object()}}; - } - std::shared_ptr video_diagnostics; - { - std::lock_guard lock(d->consumers_mutex); - if (d->consumer) - video_diagnostics = d->consumer->video_diagnostics; - } + const auto diagnostics = d->state.diagnostics.load( + std::memory_order_acquire); + nlohmann::json output = diagnostics + ? *diagnostics + : nlohmann::json{{"kind", "gallery_metrics"}, + {"protocol", "aethera.gallery.video"}, + {"version", 2}, + {"sources", nlohmann::json::object()}}; + const auto current = d->consumer.load(std::memory_order_acquire); + const auto video_diagnostics = current + ? current->video_diagnostics : nullptr; if (video_diagnostics) output["webrtc"] = (*video_diagnostics)(); return output; } diff --git a/web_server/src/Gallery_WebSocket.cpp b/web_server/src/Gallery_WebSocket.cpp index be5da24..14d5409 100644 --- a/web_server/src/Gallery_WebSocket.cpp +++ b/web_server/src/Gallery_WebSocket.cpp @@ -27,8 +27,6 @@ std::string media_failure_description(const std::exception_ptr& failure) { struct Gallery_WebSocket::Private { std::weak_ptr connection; /* 仅在连接存活时投递信令。 */ std::shared_ptr stream; /* 当前媒体组共享的编码图集。 */ - std::shared_ptr page_session; /* 当前连接所属的唯一页面权威。 */ - Page_Session::Connection_Id page_connection{}; /* 在页面权威内登记的连接标识。 */ std::unique_ptr video; /* 当前浏览器连接的 WebRTC 发送会话。 */ Gallery_Video_Stream::Stream_Id subscription{}; /* 共享图集输出的订阅标识。 */ std::atomic_bool attached{}; /* start 成功后为真,并保证 close 只执行一次。 */ @@ -37,14 +35,10 @@ struct Gallery_WebSocket::Private { Gallery_WebSocket::Gallery_WebSocket( drogon::WebSocketConnectionPtr connection, - std::shared_ptr stream, - std::shared_ptr page_session, - Page_Session::Connection_Id page_connection) + std::shared_ptr stream) : d(std::make_unique()) { d->connection = std::move(connection); d->stream = std::move(stream); - d->page_session = std::move(page_session); - d->page_connection = page_connection; } Gallery_WebSocket::~Gallery_WebSocket() { close(); } @@ -150,8 +144,6 @@ void Gallery_WebSocket::receive(std::string_view message) { void Gallery_WebSocket::close() noexcept { if (!d->attached.exchange(false, std::memory_order_acq_rel)) return; - try { d->page_session->detach(d->page_connection); } - catch (...) {} try { if (d->subscription != 0) d->stream->unsubscribe(d->subscription); } @@ -161,13 +153,11 @@ void Gallery_WebSocket::close() noexcept { } Gallery_WebSocket_Controller::Gallery_WebSocket_Controller( - Stream_Resolver value_resolver, - std::shared_ptr value_page_session) - : resolve_stream(std::move(value_resolver)), - page_session(std::move(value_page_session)) { - if (!resolve_stream || !page_session) + Stream_Resolver value_resolver) + : resolve_stream(std::move(value_resolver)) { + if (!resolve_stream) throw std::invalid_argument( - "gallery WebSocket requires media and page-session services"); + "gallery WebSocket requires a media stream resolver"); } void Gallery_WebSocket_Controller::initPathRouting() { @@ -179,43 +169,8 @@ void Gallery_WebSocket_Controller::handleNewConnection( const drogon::HttpRequestPtr& request, const drogon::WebSocketConnectionPtr& connection) { try { - const auto page_connection = page_session->attach( - request->getParameter("session"), [weak = std::weak_ptr(connection)] { - const auto active = weak.lock(); - if (!active || !active->connected()) return; - try { - active->send(nlohmann::json{ - {"kind", "page_session_replaced"}, - {"protocol", "aethera.page.session"}, - {"version", 1}}.dump(), - drogon::WebSocketMessageType::Text); - } - catch (...) {} - try { - active->shutdown(drogon::CloseCode::kViolation, - "Page session replaced"); - } - catch (...) {} - }); - if (!page_connection) { - try { - connection->send(nlohmann::json{ - {"kind", "page_session_rejected"}, - {"protocol", "aethera.page.session"}, - {"version", 1}}.dump(), - drogon::WebSocketMessageType::Text); - } - catch (...) {} - try { - connection->shutdown(drogon::CloseCode::kViolation, - "Stale page session"); - } - catch (...) {} - return; - } const auto stream = resolve_stream(request->getParameter("group")); if (!stream) { - page_session->detach(*page_connection); try { connection->send(nlohmann::json{ {"kind", "gallery_error"}, @@ -234,7 +189,7 @@ void Gallery_WebSocket_Controller::handleNewConnection( } try { auto socket = std::make_shared( - connection, stream, page_session, *page_connection); + connection, stream); connection->setContext(socket); connection->setPingMessage( "aethera-gallery-video", std::chrono::seconds(20)); @@ -243,8 +198,6 @@ void Gallery_WebSocket_Controller::handleNewConnection( catch (...) { if (const auto socket = connection->getContext()) socket->close(); - else - page_session->detach(*page_connection); try { connection->send(nlohmann::json{ {"kind", "gallery_error"}, diff --git a/web_server/src/Gallery_WebSocket.hpp b/web_server/src/Gallery_WebSocket.hpp index 6370f69..b70b276 100644 --- a/web_server/src/Gallery_WebSocket.hpp +++ b/web_server/src/Gallery_WebSocket.hpp @@ -1,6 +1,5 @@ #pragma once #include "Gallery_Video_Stream.hpp" -#include "Page_Session.hpp" #include #include #include @@ -12,9 +11,7 @@ class Gallery_WebSocket final : public std::enable_shared_from_this { public: Gallery_WebSocket(drogon::WebSocketConnectionPtr connection, - std::shared_ptr stream, - std::shared_ptr page_session, - Page_Session::Connection_Id page_connection); + std::shared_ptr stream); ~Gallery_WebSocket(); Gallery_WebSocket(const Gallery_WebSocket&) = delete; Gallery_WebSocket& operator=(const Gallery_WebSocket&) = delete; @@ -34,8 +31,7 @@ public: using Stream_Resolver = std::function(std::string_view)>; explicit Gallery_WebSocket_Controller( - Stream_Resolver resolve_stream, - std::shared_ptr page_session); + Stream_Resolver resolve_stream); void handleNewMessage(const drogon::WebSocketConnectionPtr& connection, std::string&& message, const drogon::WebSocketMessageType& type) override; @@ -45,6 +41,5 @@ public: static void initPathRouting(); private: Stream_Resolver resolve_stream; /* 媒体组标识到共享视频流的解析入口。 */ - std::shared_ptr page_session; /* 所有 Gallery 与 Plot 连接共享的页面权威。 */ }; } diff --git a/web_server/src/Graph_WebSocket.cpp b/web_server/src/Graph_WebSocket.cpp index e44d729..5d9c4e7 100644 --- a/web_server/src/Graph_WebSocket.cpp +++ b/web_server/src/Graph_WebSocket.cpp @@ -4,7 +4,6 @@ #include #include #include -#include #include #include @@ -26,24 +25,17 @@ std::uint32_t input_dimension(std::uint32_t value, std::uint32_t minimum, struct Graph_WebSocket::Private { std::weak_ptr connection; /* 仅在连接存活时投递控制消息。 */ std::shared_ptr plot; /* 该控制连接绑定的引擎与界面桥接对象。 */ - std::shared_ptr page_session; /* 当前连接所属的唯一页面权威。 */ - Page_Session::Connection_Id page_connection{}; /* 在页面权威内登记的连接标识。 */ Plot::Stream_Id stream{}; /* Plot 完成帧通知订阅标识。 */ std::atomic_bool attached{}; /* start 成功后为真,并保证 close 只执行一次。 */ - std::mutex viewport_mutex; /* 输入解码读取视口尺寸的短临界区。 */ - std::uint32_t width{320}; /* 当前图集槽位对应的输入坐标宽度。 */ - std::uint32_t height{192}; /* 当前图集槽位对应的输入坐标高度。 */ + /* 宽高属于一个不可拆分的输入坐标版本;单个 64 位原子值是唯一状态源。 */ + std::atomic_uint64_t viewport{(std::uint64_t{320} << 32U) | 192U}; }; Graph_WebSocket::Graph_WebSocket(drogon::WebSocketConnectionPtr connection, - std::shared_ptr plot, - std::shared_ptr page_session, - Page_Session::Connection_Id page_connection) + std::shared_ptr plot) : d(std::make_unique()) { d->connection = std::move(connection); d->plot = std::move(plot); - d->page_session = std::move(page_session); - d->page_connection = page_connection; } Graph_WebSocket::~Graph_WebSocket() { close(); } @@ -85,9 +77,9 @@ void Graph_WebSocket::receive(std::string_view message) { width = input_dimension(viewport->value("width", 320U), 160U, 1920U); height = input_dimension(viewport->value("height", 192U), 120U, 1080U); } - std::lock_guard lock(d->viewport_mutex); - d->width = width; - d->height = height; + d->viewport.store((static_cast(width) << 32U) | + static_cast(height), + std::memory_order_release); return; } if (kind == "manual_render") { @@ -101,13 +93,9 @@ void Graph_WebSocket::receive(std::string_view message) { const auto type = magic_enum::enum_cast( input->value("type", std::string{})); if (!type) return; - std::uint32_t width{}; - std::uint32_t height{}; - { - std::lock_guard lock(d->viewport_mutex); - width = d->width; - height = d->height; - } + const auto viewport = d->viewport.load(std::memory_order_acquire); + const auto width = static_cast(viewport >> 32U); + const auto height = static_cast(viewport); Plot_Input_Event decoded; decoded.type = *type; const auto read_point = [&](std::string_view key, @@ -152,20 +140,14 @@ void Graph_WebSocket::receive(std::string_view message) { void Graph_WebSocket::close() noexcept { if (!d->attached.exchange(false, std::memory_order_acq_rel)) return; - try { d->page_session->detach(d->page_connection); } - catch (...) {} try { d->plot->unsubscribe(d->stream); } catch (...) {} } -Graph_WebSocket_Controller::Graph_WebSocket_Controller( - Plot_Resolver resolver, - std::shared_ptr value_page_session) - : resolve_plot(std::move(resolver)), - page_session(std::move(value_page_session)) { - if (!resolve_plot || !page_session) - throw std::invalid_argument( - "plot WebSocket requires plot and page-session services"); +Graph_WebSocket_Controller::Graph_WebSocket_Controller(Plot_Resolver resolver) + : resolve_plot(std::move(resolver)) { + if (!resolve_plot) + throw std::invalid_argument("plot WebSocket requires a plot resolver"); } void Graph_WebSocket_Controller::initPathRouting() { @@ -177,43 +159,8 @@ void Graph_WebSocket_Controller::handleNewConnection( const drogon::HttpRequestPtr& request, const drogon::WebSocketConnectionPtr& connection) { try { - const auto page_connection = page_session->attach( - request->getParameter("session"), [weak = std::weak_ptr(connection)] { - const auto active = weak.lock(); - if (!active || !active->connected()) return; - try { - active->send(nlohmann::json{ - {"kind", "page_session_replaced"}, - {"protocol", "aethera.page.session"}, - {"version", 1}}.dump(), - drogon::WebSocketMessageType::Text); - } - catch (...) {} - try { - active->shutdown(drogon::CloseCode::kViolation, - "Page session replaced"); - } - catch (...) {} - }); - if (!page_connection) { - try { - connection->send(nlohmann::json{ - {"kind", "page_session_rejected"}, - {"protocol", "aethera.page.session"}, - {"version", 1}}.dump(), - drogon::WebSocketMessageType::Text); - } - catch (...) {} - try { - connection->shutdown(drogon::CloseCode::kViolation, - "Stale page session"); - } - catch (...) {} - return; - } auto plot = resolve_plot(graph_id_from_path(request->path())); if (!plot) { - page_session->detach(*page_connection); try { connection->shutdown(drogon::CloseCode::kViolation, "Unknown Aethera plot"); @@ -223,7 +170,7 @@ void Graph_WebSocket_Controller::handleNewConnection( } try { auto socket = std::make_shared( - connection, std::move(plot), page_session, *page_connection); + connection, std::move(plot)); connection->setContext(socket); connection->setPingMessage( "aethera-gallery-plot", std::chrono::seconds(20)); @@ -232,8 +179,6 @@ void Graph_WebSocket_Controller::handleNewConnection( catch (...) { if (const auto socket = connection->getContext()) socket->close(); - else - page_session->detach(*page_connection); try { connection->shutdown(drogon::CloseCode::kViolation, "Aethera plot setup failed"); diff --git a/web_server/src/Graph_WebSocket.hpp b/web_server/src/Graph_WebSocket.hpp index 11faa68..c34827f 100644 --- a/web_server/src/Graph_WebSocket.hpp +++ b/web_server/src/Graph_WebSocket.hpp @@ -1,5 +1,4 @@ #pragma once -#include "Page_Session.hpp" #include "Plot.hpp" #include #include @@ -11,9 +10,7 @@ namespace aethera::web { class Graph_WebSocket final : public std::enable_shared_from_this { public: Graph_WebSocket(drogon::WebSocketConnectionPtr connection, - std::shared_ptr plot, - std::shared_ptr page_session, - Page_Session::Connection_Id page_connection); + std::shared_ptr plot); ~Graph_WebSocket(); Graph_WebSocket(const Graph_WebSocket&) = delete; Graph_WebSocket& operator=(const Graph_WebSocket&) = delete; @@ -30,9 +27,7 @@ class Graph_WebSocket_Controller final : public drogon::WebSocketController { public: using Plot_Resolver = std::function(std::string_view)>; - explicit Graph_WebSocket_Controller( - Plot_Resolver resolver, - std::shared_ptr page_session); + explicit Graph_WebSocket_Controller(Plot_Resolver resolver); void handleNewMessage(const drogon::WebSocketConnectionPtr& connection, std::string&& message, const drogon::WebSocketMessageType& type) override; @@ -42,6 +37,5 @@ public: static void initPathRouting(); private: Plot_Resolver resolve_plot; /* 图标识到引擎 Plot 的解析入口。 */ - std::shared_ptr page_session; /* 所有 Gallery 与 Plot 连接共享的页面权威。 */ }; } diff --git a/web_server/src/Page_Session.cpp b/web_server/src/Page_Session.cpp deleted file mode 100644 index 6ff256c..0000000 --- a/web_server/src/Page_Session.cpp +++ /dev/null @@ -1,80 +0,0 @@ -#include "Page_Session.hpp" -#include -#include -#include -#include -#include -#include -#include -#include -#include - -namespace aethera::web { -namespace { -std::string make_page_token() { - std::random_device random; - std::array words{}; - for (auto& word : words) word = random(); - std::ostringstream token; - token << std::hex << std::setfill('0'); - for (const auto word : words) token << std::setw(8) << word; - return std::move(token).str(); -} -} - -struct Page_Session::Private { - std::mutex mutex; /* 权威页面切换与连接登记的唯一临界区。 */ - std::string active_token; /* 当前唯一可建立连接的页面令牌。 */ - Connection_Id next_connection{1}; /* 进程生命周期内递增的连接登记标识。 */ - std::unordered_map - connections; /* 当前页面拥有的连接关闭动作。 */ -}; - -Page_Session::Page_Session() : d(std::make_unique()) {} - -Page_Session::~Page_Session() = default; - -std::string Page_Session::begin_page() { - auto token = make_page_token(); - std::vector displaced; - { - std::lock_guard lock(d->mutex); - while (token == d->active_token) token = make_page_token(); - displaced.reserve(d->connections.size()); - for (auto& [connection, close_connection] : d->connections) { - static_cast(connection); - displaced.push_back(std::move(close_connection)); - } - d->connections.clear(); - d->active_token = token; - } - - /* 页面接管只在临界区内切换权威令牌。关闭第三方连接必须在锁外执行, - * 因而 WebSocket 关闭回调可以安全地反向调用 detach(),且不会阻塞新页登记。 */ - for (auto& close_connection : displaced) { - try { close_connection(); } - catch (...) { - /* 单条旧连接已经失去权威资格;其关闭失败不能妨碍其余连接被隔离。 */ - } - } - return token; -} - -std::expected -Page_Session::attach(std::string_view token, - Close_Connection close_connection) { - if (!close_connection) - throw std::invalid_argument("page connection requires a close action"); - std::lock_guard lock(d->mutex); - if (token.empty() || token != d->active_token) - return std::unexpected(Attach_Result::stale_session); - const auto connection = d->next_connection++; - d->connections.emplace(connection, std::move(close_connection)); - return connection; -} - -void Page_Session::detach(Connection_Id connection) { - std::lock_guard lock(d->mutex); - d->connections.erase(connection); -} -} diff --git a/web_server/src/Page_Session.hpp b/web_server/src/Page_Session.hpp deleted file mode 100644 index 6807b58..0000000 --- a/web_server/src/Page_Session.hpp +++ /dev/null @@ -1,33 +0,0 @@ -#pragma once -#include -#include -#include -#include -#include -#include - -namespace aethera::web { -class Page_Session final { -public: - using Connection_Id = std::uint64_t; - using Close_Connection = std::function; - - enum class Attach_Result : std::uint8_t { - stale_session - }; - - Page_Session(); - ~Page_Session(); - Page_Session(const Page_Session&) = delete; - Page_Session& operator=(const Page_Session&) = delete; - - [[nodiscard]] std::string begin_page(); - [[nodiscard]] std::expected attach( - std::string_view token, Close_Connection close_connection); - void detach(Connection_Id connection); - -private: - struct Private; - std::unique_ptr d; -}; -} diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index c6c512b..3a953f7 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -16,9 +16,7 @@ #include #include #include -#include #include -#include #include #include #include @@ -64,16 +62,15 @@ struct Frame_Pacing_Properties { class Frame_Policy final { public: - [[nodiscard]] Frame_Pacing_Properties snapshot() const { - std::lock_guard lock(mutex); - return pacing; + [[nodiscard]] Frame_Pacing_Properties read() const { + return *pacing.load(std::memory_order_acquire); } [[nodiscard]] nlohmann::json schema() const; [[nodiscard]] nlohmann::json write_prop(std::string_view key, const nlohmann::json& value); private: - mutable std::mutex mutex; - Frame_Pacing_Properties pacing{}; + std::atomic> pacing{ + std::make_shared()}; }; std::string_view pacing_mode_name(Frame_Pacing_Mode mode) { @@ -101,22 +98,22 @@ std::string_view pixel_format_name(render_3d::Pixel_Format format) { } nlohmann::json Frame_Policy::schema() const { - std::lock_guard lock(mutex); + const auto current = read(); return { {"id", "frame-analysis"}, {"label", "渲染与媒体流水线"}, {"kind", "analysis"}, {"fields", nlohmann::json::array({ {{"key", "render_enabled"}, {"label", "持续渲染与采样"}, {"editor", "boolean"}, {"editable", true}, {"description", "控制共享帧时钟是否继续调用当前 Scene;画面隐藏不会修改此项。"}, {"technical_description", "Authoritative server-side render and sampling switch."}, - {"value", pacing.render_enabled}}, + {"value", current.render_enabled}}, {{"key", "video_enabled"}, {"label", "图集视频传输"}, {"editor", "boolean"}, {"editable", true}, {"description", "控制完成帧是否进入页面级 RGBA 图集;默认开启,用于完整链路压测。"}, {"technical_description", "Authoritative tile publication switch for the shared gallery video."}, - {"value", pacing.video_enabled}}, + {"value", current.video_enabled}}, {{"key", "pacing_mode"}, {"label", "服务端帧策略"}, {"editor", "select"}, {"editable", true}, {"description", "只控制 Scene::render 的调用节奏;完成回调只负责归还帧并发布结果。"}, {"technical_description", "Render policy driven by the common 100 Hz gallery clock."}, - {"value", pacing_mode_name(pacing.mode)}, + {"value", pacing_mode_name(current.mode)}, {"options", nlohmann::json::array({ {{"value", "manual"}, {"label", "手动渲染"}}, {{"value", "fixed_rate"}, {"label", "固定频率"}}, @@ -126,19 +123,33 @@ nlohmann::json Frame_Policy::schema() const { {"editable", true}, {"minimum", 0.1}, {"maximum", 100.0}, {"step", 0.1}, {"description", "固定频率模式下从页面级 100 Hz 时钟采样当前 Scene 的次数。"}, {"technical_description", "Per-plot rate selected from the common gallery render timeline."}, - {"value", pacing.fixed_rate_fps}} + {"value", current.fixed_rate_fps}} })} }; } nlohmann::json Frame_Policy::write_prop(std::string_view key, const nlohmann::json& value) { - std::lock_guard lock(mutex); + const auto update = [this](auto&& edit) { + auto current = pacing.load(std::memory_order_acquire); + for (;;) { + auto next = std::make_shared(*current); + edit(*next); + std::shared_ptr desired = next; + if (pacing.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + return desired; + } + }; if (key == "render_enabled" || key == "video_enabled") { if (!value.is_boolean()) return {{"success", false}, {"error", "frame policy switch requires a boolean"}}; - bool& target = key == "render_enabled" ? pacing.render_enabled : pacing.video_enabled; - target = value.get(); + const bool target = value.get(); + update([&](Frame_Pacing_Properties& next) { + (key == "render_enabled" ? next.render_enabled : next.video_enabled) = + target; + }); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, {"value", target}}; } @@ -148,9 +159,9 @@ nlohmann::json Frame_Policy::write_prop(std::string_view key, const auto parsed = parse_pacing_mode(value.get_ref()); if (!parsed) return {{"success", false}, {"error", "unknown frame pacing mode"}}; - pacing.mode = *parsed; + update([&](Frame_Pacing_Properties& next) { next.mode = *parsed; }); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, - {"value", pacing_mode_name(pacing.mode)}}; + {"value", pacing_mode_name(*parsed)}}; } if (key == "fixed_rate_fps") { if (!value.is_number()) @@ -158,9 +169,11 @@ nlohmann::json Frame_Policy::write_prop(std::string_view key, const double next = value.get(); if (!std::isfinite(next) || next < 0.1 || next > 100.0) return {{"success", false}, {"error", "fixed_rate_fps must be between 0.1 and 100"}}; - pacing.fixed_rate_fps = next; + update([&](Frame_Pacing_Properties& properties) { + properties.fixed_rate_fps = next; + }); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, - {"value", pacing.fixed_rate_fps}}; + {"value", next}}; } return {{"success", false}, {"error", "unknown frame runtime property"}}; } @@ -332,7 +345,7 @@ struct Plot::Private { struct Managed_Frame { std::chrono::microseconds presentation_time{}; /* 共享页面时钟产生的媒体时间戳。 */ Frame frame{}; /* 三缓冲物理槽拥有且反复承载逻辑帧。 */ - Frame_State state{Frame_State::available}; /* 本槽唯一生命周期状态。 */ + std::atomic state{Frame_State::available}; /* 本槽唯一生命周期状态。 */ }; struct Consumer { @@ -340,29 +353,28 @@ struct Plot::Private { std::uint32_t width{}; std::uint32_t height{}; }; + using Consumer_Map = std::unordered_map; struct Stream_Snapshot { - std::vector> consumers; + std::shared_ptr consumers; std::uint32_t width{}; std::uint32_t height{}; }; std::unique_ptr view; std::once_flag start_once; - mutable std::mutex consumers_mutex; - std::unordered_map consumers; /* 媒体图集和诊断连接的唯一订阅表。 */ + std::atomic> consumers{ + std::make_shared()}; /* 低频订阅修改发布不可变版本。 */ std::atomic_uint64_t next_stream_id{1}; std::atomic> terminal_failure{}; /* 首次 Plot Unknown Failure 的唯一终止状态。 */ std::uint64_t next_frame_sequence{1}; Frame_Policy frame_policy{}; - mutable std::mutex frame_mutex; static constexpr std::size_t scene_frame_capacity{3}; std::array frame_slots{}; /* Scene 借用的稳定三缓冲物理帧。 */ - std::optional callback_retired_slot{}; /* 上一个回调返回后才可重新开始的槽。 */ + std::atomic_size_t callback_retired_slot{scene_frame_capacity}; Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */ - mutable std::mutex tick_mutex; - std::optional pending_tick{}; /* 时钟拥塞时只保留尚未处理的最新时间点。 */ - bool tick_task_scheduled{}; /* Taskflow 中是否已有唯一 tick 消费任务。 */ + std::atomic> pending_tick{}; + std::atomic_bool tick_task_scheduled{}; /* 唯一短任务准入;不占用 Worker 等待。 */ double last_clock_render_time_ms{-std::numeric_limits::infinity()}; std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()}; std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */ @@ -373,9 +385,12 @@ struct Plot::Private { std::atomic_uint64_t scene_rejection_count{}; /* Scene 单帧准入拒绝的提交次数。 */ std::atomic_uint64_t submitted_frame_count{}; /* 成功提交给 Scene 的帧总数。 */ std::atomic_size_t taskflow_trace_remaining{}; /* 尚待标记的实际渲染帧数。 */ - mutable std::mutex taskflow_trace_mutex{}; /* 只保护低频请求结果的交换。 */ - std::size_t taskflow_trace_requested_count{}; /* 当前批次请求总帧数。 */ - std::vector taskflow_traces{}; /* 已完成帧直接发布的 DAG 与 Observer 结果。 */ + static constexpr std::size_t maximum_taskflow_trace_frames{120}; + /* 高 32 位 requested,低 32 位 captured。每槽只发布一次不可变 Trace, + * GET 直接读取已发布槽位,不复制或重排整个历史容器。 */ + 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 对子图的借用。 */ @@ -444,10 +459,9 @@ nlohmann::json Plot::Private::schema() const { Plot::Private::Stream_Snapshot Plot::Private::stream_snapshot() const { Stream_Snapshot result; - std::lock_guard lock(consumers_mutex); - result.consumers.reserve(consumers.size()); - for (const auto& [id, consumer] : consumers) { - result.consumers.emplace_back(id, consumer); + result.consumers = consumers.load(std::memory_order_acquire); + for (const auto& [id, consumer] : *result.consumers) { + static_cast(id); if (consumer.width == 0 || consumer.height == 0) continue; result.width = std::max(result.width, consumer.width); result.height = std::max(result.height, consumer.height); @@ -463,7 +477,7 @@ void Plot::Private::publish( try { const auto snapshot = stream_snapshot(); std::vector failed_consumers; - for (const auto& [id, consumer] : snapshot.consumers) { + for (const auto& [id, consumer] : *snapshot.consumers) { if (!consumer.handler) continue; try { consumer.handler(frame); @@ -473,27 +487,27 @@ void Plot::Private::publish( } } if (failed_consumers.empty()) return; - std::lock_guard lock(consumers_mutex); - for (const auto id : failed_consumers) consumers.erase(id); + auto current = consumers.load(std::memory_order_acquire); + for (;;) { + auto next = std::make_shared(*current); + for (const auto id : failed_consumers) next->erase(id); + std::shared_ptr desired = next; + if (consumers.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + break; + } } catch (...) {} } void Plot::Private::consume_tick(std::weak_ptr lifetime) { - std::optional tick; - { - std::lock_guard lock(tick_mutex); - tick = std::exchange(pending_tick, {}); - } + const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel); if (tick) clock_tick(*tick); - bool schedule_again{}; - { - std::lock_guard lock(tick_mutex); - schedule_again = pending_tick.has_value(); - if (!schedule_again) tick_task_scheduled = false; - } - if (schedule_again) { + 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; @@ -505,7 +519,7 @@ void Plot::Private::consume_tick(std::weak_ptr lifetime) { void Plot::Private::clock_tick(const Plot_Render_Tick& tick) { if (terminal_failure.load(std::memory_order_acquire)) return; - const auto pacing = frame_policy.snapshot(); + const auto pacing = frame_policy.read(); if (!pacing.render_enabled || pacing.mode == Frame_Pacing_Mode::manual) { policy_skip_count.fetch_add(1, std::memory_order_relaxed); return; @@ -539,46 +553,43 @@ bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) { 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.snapshot(); - if (!pacing.render_enabled || streams.consumers.empty()) return; + const auto pacing = frame_policy.read(); + if (!pacing.render_enabled || streams.consumers->empty()) return; std::size_t slot_index{}; Managed_Frame* managed{}; - { - std::lock_guard lock(frame_mutex); - /* - * Frame_State 是 Plot 数据准备生命周期的唯一权威来源。必须在修改 - * Scene/Visual 输入之前取得准入;Scene::render() 内部再拒绝已经太晚, - * 因为上一帧的异步 prepare 可能正在读取同一份业务数据。 - */ - if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) { - return slot.state == Frame_State::in_flight; - })) { - preparation_busy_count.fetch_add(1, std::memory_order_relaxed); - return; - } - const auto available = std::ranges::find_if( - frame_slots, [](const Managed_Frame& slot) { - return slot.state == Frame_State::available; - }); - if (available == frame_slots.end()) { - frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed); - return; - } - slot_index = static_cast( - std::distance(frame_slots.begin(), available)); - managed = &*available; - managed->state = Frame_State::in_flight; - managed->presentation_time = - std::chrono::duration_cast( - std::chrono::duration( - tick.time_milliseconds)); + /* 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; + })) { + preparation_busy_count.fetch_add(1, std::memory_order_relaxed); + return; } + 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, + std::memory_order_acq_rel, std::memory_order_acquire)) + continue; + slot_index = index; + managed = &frame_slots[index]; + break; + } + if (!managed) { + frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed); + return; + } + managed->presentation_time = + std::chrono::duration_cast( + std::chrono::duration(tick.time_milliseconds)); const auto rollback_unsubmitted = [this, slot_index] { - std::lock_guard lock(frame_mutex); auto& slot = frame_slots[slot_index]; - if (slot.state == Frame_State::in_flight) - slot.state = Frame_State::available; + auto expected = Frame_State::in_flight; + static_cast(slot.state.compare_exchange_strong( + expected, Frame_State::available, std::memory_order_acq_rel, + std::memory_order_acquire)); }; bool taskflow_trace_claimed{}; const auto restore_taskflow_trace_claim = [this, &taskflow_trace_claimed] { @@ -662,36 +673,32 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { void Plot::Private::publish_completed_frame() { Render_Frame* frame{}; Managed_Frame* managed{}; - { - std::lock_guard lock(frame_mutex); - /* - * Scene 在用户回调返回之后才标记 callback_finished,因此当前回调 - * 不能立即重启自己的物理帧。下一个串行完成回调开始时,上一个 - * callback_retired 槽已经确定离开 Scene,可安全归还。三槽由此在 - * 不改变 Scene 回调契约的前提下形成稳定的循环生命周期。 - */ - if (callback_retired_slot) { - auto& retired = frame_slots[*callback_retired_slot]; - if (retired.state != Frame_State::callback_retired) - throw std::logic_error("Plot retired frame state is inconsistent"); - retired.state = Frame_State::available; - callback_retired_slot.reset(); - } - for (std::size_t index = 0; index < frame_slots.size(); ++index) { - if (frame_slots[index].state != Frame_State::in_flight) continue; - if (managed) - throw std::logic_error("Plot has multiple active Scene frames"); - frame = std::visit( - [](const auto& value) -> Render_Frame* { return value.get(); }, - frame_slots[index].frame); - managed = &frame_slots[index]; - } - if (!managed || managed->state != Frame_State::in_flight) - throw std::logic_error("Scene completion graph has no active Plot frame"); + /* 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; + if (managed) + throw std::logic_error("Plot has multiple active Scene frames"); + 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"); try { - const auto pacing = frame_policy.snapshot(); + const auto pacing = frame_policy.read(); const auto identity = frame->identity(); Frame_Identity rendered_identity = identity; std::shared_ptr> pixel_storage; @@ -748,30 +755,44 @@ void Plot::Private::retire_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}; - { - std::lock_guard lock(frame_mutex); - for (std::size_t index = 0; index < frame_slots.size(); ++index) { - auto* address = std::visit( - [](const auto& value) -> Render_Frame* { return value.get(); }, - frame_slots[index].frame); - if (address != frame) continue; - if (frame_slots[index].state != Frame_State::in_flight) - throw std::logic_error("completed Plot frame is not in flight"); - slot_index = index; - break; - } - if (slot_index == scene_frame_capacity) - throw std::logic_error( - "frame callback has no externally owned active frame"); - if (callback_retired_slot) - throw std::logic_error("Plot has more than one callback-retired frame"); - frame_slots[slot_index].state = Frame_State::callback_retired; - callback_retired_slot = slot_index; + for (std::size_t index = 0; index < frame_slots.size(); ++index) { + auto* address = std::visit( + [](const auto& value) -> Render_Frame* { return value.get(); }, + frame_slots[index].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; + 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, + std::memory_order_acquire)) + throw std::logic_error("Plot has more than one callback-retired frame"); if (frame->taskflow_trace_requested()) { - auto trace = frame->taskflow_trace(); - std::lock_guard lock(taskflow_trace_mutex); - taskflow_traces.push_back(std::move(trace)); + 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; + } } } @@ -837,14 +858,18 @@ Plot::Stream_Id Plot::subscribe(Stream_Handler handler) { throw std::invalid_argument("Plot subscription requires a handler"); ensure_started(); const auto id = d->next_stream_id.fetch_add(1, std::memory_order_relaxed); - Stream_Handler notification; - { - std::lock_guard lock(d->consumers_mutex); - d->consumers.emplace(id, Private::Consumer{std::move(handler)}); - if (d->terminal_failure.load(std::memory_order_acquire)) - notification = d->consumers.at(id).handler; + const auto notification = handler; + auto current = d->consumers.load(std::memory_order_acquire); + for (;;) { + auto next = std::make_shared(*current); + next->emplace(id, Private::Consumer{handler}); + std::shared_ptr desired = next; + if (d->consumers.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + break; } - if (notification) { + if (d->terminal_failure.load(std::memory_order_acquire)) { const auto failure = d->terminal_failure.load(std::memory_order_acquire); try { notification(std::make_shared( @@ -863,35 +888,44 @@ Plot::Stream_Id Plot::subscribe(Stream_Handler handler) { } void Plot::unsubscribe(Stream_Id stream) { - std::lock_guard lock(d->consumers_mutex); - d->consumers.erase(stream); + auto current = d->consumers.load(std::memory_order_acquire); + while (current->contains(stream)) { + auto next = std::make_shared(*current); + next->erase(stream); + std::shared_ptr desired = next; + if (d->consumers.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + return; + } } void Plot::configure_stream(Stream_Id stream, std::uint32_t width, std::uint32_t height) { - std::lock_guard lock(d->consumers_mutex); - const auto found = d->consumers.find(stream); - if (found == d->consumers.end()) return; - found->second.width = width; - found->second.height = height; + auto current = d->consumers.load(std::memory_order_acquire); + for (;;) { + const auto found = current->find(stream); + if (found == current->end()) return; + auto next = std::make_shared(*current); + auto& consumer = next->at(stream); + consumer.width = width; + consumer.height = height; + std::shared_ptr desired = next; + if (d->consumers.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + return; + } } 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); - bool schedule{}; - { - std::lock_guard lock(d->tick_mutex); - if (d->pending_tick) - d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed); - d->pending_tick = tick; - if (!d->tick_task_scheduled) { - d->tick_task_scheduled = true; - schedule = true; - } - } - if (!schedule) return; + 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(); @@ -995,7 +1029,7 @@ nlohmann::json Plot::diagnostics() const { } }, d->scene); - const auto pacing = d->frame_policy.snapshot(); + const auto pacing = d->frame_policy.read(); const auto stream = d->stream_snapshot(); nlohmann::json supported_formats = nlohmann::json::array(); if (is_3d) { @@ -1068,7 +1102,6 @@ nlohmann::json Plot::diagnostics() const { {"callback_max_ms", milliseconds(gpu.callback_max_ns)}, {"callback_failure_count", gpu.callback_failure_count}, {"backpressure_count", gpu.backpressure_count}, - {"backpressure_wait_ms", milliseconds(gpu.backpressure_wait_ns)}, {"fault_count", gpu.fault_count}, {"abandoned_count", gpu.abandoned_count}, {"stopping", gpu.stopping}}; @@ -1079,29 +1112,37 @@ nlohmann::json Plot::diagnostics() const { } void Plot::request_taskflow_trace(std::size_t frame_count) { - if (frame_count == 0 || frame_count > 120) + if (frame_count == 0 || + frame_count > Private::maximum_taskflow_trace_frames) throw std::invalid_argument("Taskflow trace frame_count must be between 1 and 120"); ensure_started(); - { - std::lock_guard lock(d->taskflow_trace_mutex); - if (d->taskflow_trace_requested_count != d->taskflow_traces.size()) + auto control = d->taskflow_trace_control.load(std::memory_order_acquire); + for (;;) { + const auto requested = static_cast(control >> 32U); + const auto captured = static_cast(control); + if (requested != captured) throw std::logic_error("A Taskflow frame trace request is already active"); - d->taskflow_traces.clear(); - d->taskflow_traces.reserve(frame_count); - d->taskflow_trace_requested_count = frame_count; + const auto next = static_cast(frame_count) << 32U; + if (d->taskflow_trace_control.compare_exchange_weak( + control, next, std::memory_order_release, + std::memory_order_acquire)) + break; } + for (auto& slot : d->taskflow_trace_slots) + slot.store({}, std::memory_order_release); d->taskflow_trace_remaining.store(frame_count, std::memory_order_release); } nlohmann::json Plot::taskflow_trace() const { nlohmann::json frames = nlohmann::json::array(); - std::size_t requested{}; - { - std::lock_guard lock(d->taskflow_trace_mutex); - requested = d->taskflow_trace_requested_count; - for (const auto& trace : d->taskflow_traces) - frames.push_back(taskflow_trace_json(trace)); - } + 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( + std::memory_order_acquire)) + frames.push_back(taskflow_trace_json(*trace)); const auto remaining = d->taskflow_trace_remaining.load( std::memory_order_acquire); return { diff --git a/web_server/src/WebRtc_Video_Session.cpp b/web_server/src/WebRtc_Video_Session.cpp index d273bd7..841738e 100644 --- a/web_server/src/WebRtc_Video_Session.cpp +++ b/web_server/src/WebRtc_Video_Session.cpp @@ -1,6 +1,5 @@ #include "WebRtc_Video_Session.hpp" #include -#include #include #include #include @@ -37,32 +36,6 @@ void update_maximum(std::atomic_uint64_t& maximum, } struct WebRtc_Video_Session::Private { - struct State_Tag {}; - struct State : double_buffer::State_Type { - std::size_t outstanding_video_frames{}; - std::size_t transport_buffered_bytes{}; - std::uint64_t queued_frame_count{}; - std::uint64_t rejected_frame_count{}; - std::uint64_t queued_byte_count{}; - std::uint64_t sent_frame_count{}; - std::uint64_t sent_byte_count{}; - std::uint64_t send_failure_count{}; - std::uint64_t send_total_ns{}; - std::uint64_t send_max_ns{}; - std::uint64_t current_send_ns{}; - std::uint64_t media_lock_wait_count{}; - std::uint64_t media_lock_wait_total_ns{}; - std::uint64_t media_lock_wait_max_ns{}; - std::uint64_t readiness_check_count{}; - std::uint64_t readiness_reject_count{}; - std::uint64_t close_join_total_ns{}; - std::uint64_t close_join_max_ns{}; - std::uint64_t current_close_join_ns{}; - bool track_open{}; - bool closed{}; - bool sender_running{}; - bool operator==(const State&) const = default; - }; struct Callback_State { Signal_Handler signal_handler; /* SDP/ICE 信令的唯一交付出口。 */ Ready_Handler ready_handler; /* Track 打开后请求首个 IDR。 */ @@ -82,9 +55,11 @@ struct WebRtc_Video_Session::Private { std::shared_ptr callbacks; std::shared_ptr peer; - std::shared_ptr video_track; + std::atomic> video_track{}; std::shared_ptr rtp_config; - std::mutex media_mutex; /* 发送、信令修改与关闭的生命周期边界。 */ + /* libdatachannel 没有为同一 Peer 的 SDP/ICE 修改与 close 提供可并发契约; + * 这里只串行化第三方对象生命周期。视频热路径通过原子 Track 发布,不取锁。 */ + std::mutex media_mutex; moodycamel::BlockingConcurrentQueue pending_frames{4}; /* 尚未发送的最新编码帧入口。 */ std::thread sender_thread{}; /* 唯一的 RTP/SRTP 发送资源域。 */ @@ -100,9 +75,6 @@ struct WebRtc_Video_Session::Private { std::atomic_uint64_t send_started_ns{}; std::atomic_bool send_active{}; std::atomic_bool sender_loop_active{}; - std::atomic_uint64_t media_lock_wait_count{}; - std::atomic_uint64_t media_lock_wait_total_ns{}; - std::atomic_uint64_t media_lock_wait_max_ns{}; std::atomic_size_t transport_buffered_bytes{}; std::atomic_uint64_t readiness_check_count{}; std::atomic_uint64_t readiness_reject_count{}; @@ -110,18 +82,6 @@ struct WebRtc_Video_Session::Private { std::atomic_uint64_t close_join_total_ns{}; std::atomic_uint64_t close_join_max_ns{}; std::atomic_bool close_join_active{}; - /* WebRTC 是 web_server 会话域,不借用 Scene State。热路径只写原子计数; - * 低频 diagnostics 请求在这里聚合并交换会话自己的发布双缓冲。 */ - mutable std::mutex state_publication_mutex; - mutable double_buffer::Publish_Double_Buffer state{}; - - void record_media_lock_wait(std::uint64_t started) noexcept { - const auto elapsed = steady_nanoseconds() - started; - media_lock_wait_count.fetch_add(1, std::memory_order_relaxed); - media_lock_wait_total_ns.fetch_add(elapsed, std::memory_order_relaxed); - update_maximum(media_lock_wait_max_ns, elapsed); - } - Private(Signal_Handler signal_handler, Ready_Handler ready_handler, Failure_Handler failure_handler) : callbacks(std::make_shared()) { @@ -146,20 +106,14 @@ struct WebRtc_Video_Session::Private { callbacks->closed.load(std::memory_order_acquire)) return; try { - const auto lock_started = steady_nanoseconds(); - std::shared_ptr track; - { - std::lock_guard lock(media_mutex); - record_media_lock_wait(lock_started); - if (!video_track || - !callbacks->track_open.load(std::memory_order_acquire) || - callbacks->closed.load(std::memory_order_acquire)) { - outstanding_video_frames.fetch_sub( - 1, std::memory_order_release); - rejected_frame_count.fetch_add(1, std::memory_order_relaxed); - continue; - } - track = video_track; + const auto track = video_track.load(std::memory_order_acquire); + if (!track || + !callbacks->track_open.load(std::memory_order_acquire) || + callbacks->closed.load(std::memory_order_acquire)) { + outstanding_video_frames.fetch_sub( + 1, std::memory_order_release); + rejected_frame_count.fetch_add(1, std::memory_order_relaxed); + continue; } /* rtc::binary 与编码器 access unit 使用相同的 byte vector。 * 这里只在锁内取得 Track 的共享所有权;packetizer 和网络发送均为 @@ -298,7 +252,7 @@ void WebRtc_Video_Session::start() { peer->setLocalDescription(); d->peer = std::move(peer); - d->video_track = std::move(video_track); + d->video_track.store(std::move(video_track), std::memory_order_release); d->rtp_config = std::move(rtp_config); d->sender_thread = std::thread([data = d.get()] { data->run_sender(); }); } @@ -325,7 +279,7 @@ WebRtc_Video_Session::Send_Result WebRtc_Video_Session::send( if (frame.annex_b.empty() || !d->callbacks->track_open.load(std::memory_order_acquire) || d->callbacks->closed.load(std::memory_order_acquire) || - !d->sender_thread.joinable()) { + !d->sender_loop_active.load(std::memory_order_acquire)) { d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed); return Send_Result::not_open; } @@ -362,18 +316,12 @@ WebRtc_Video_Session::Send_Result WebRtc_Video_Session::send( bool WebRtc_Video_Session::can_accept_video() const noexcept { d->readiness_check_count.fetch_add(1, std::memory_order_relaxed); - const auto lock_started = steady_nanoseconds(); - std::shared_ptr track; - { - std::lock_guard lock(d->media_mutex); - d->record_media_lock_wait(lock_started); - track = d->video_track; - } + const auto track = d->video_track.load(std::memory_order_acquire); const auto buffered = track ? track->bufferedAmount() : 0; d->transport_buffered_bytes.store(buffered, std::memory_order_relaxed); const bool ready = d->callbacks->track_open.load(std::memory_order_acquire) && !d->callbacks->closed.load(std::memory_order_acquire) && - d->sender_thread.joinable() && + d->sender_loop_active.load(std::memory_order_acquire) && track && buffered < maximum_transport_buffered_bytes && d->outstanding_video_frames.load(std::memory_order_acquire) < @@ -384,38 +332,18 @@ bool WebRtc_Video_Session::can_accept_video() const noexcept { } nlohmann::json WebRtc_Video_Session::diagnostics() const { - std::lock_guard publication_lock(d->state_publication_mutex); - auto& state_to_publish = *d->state.current; const auto now = steady_nanoseconds(); - state_to_publish.outstanding_video_frames = + const auto outstanding_video_frames = d->outstanding_video_frames.load(std::memory_order_relaxed); - state_to_publish.transport_buffered_bytes = - d->transport_buffered_bytes.load(std::memory_order_relaxed); - state_to_publish.queued_frame_count = d->queued_frame_count.load(std::memory_order_relaxed); - state_to_publish.rejected_frame_count = d->rejected_frame_count.load(std::memory_order_relaxed); - state_to_publish.queued_byte_count = d->queued_byte_count.load(std::memory_order_relaxed); - state_to_publish.sent_frame_count = d->sent_frame_count.load(std::memory_order_relaxed); - state_to_publish.sent_byte_count = d->sent_byte_count.load(std::memory_order_relaxed); - state_to_publish.send_failure_count = d->send_failure_count.load(std::memory_order_relaxed); - state_to_publish.send_total_ns = d->send_total_ns.load(std::memory_order_relaxed); - state_to_publish.send_max_ns = d->send_max_ns.load(std::memory_order_relaxed); - state_to_publish.current_send_ns = d->send_active.load(std::memory_order_acquire) - ? now - d->send_started_ns.load(std::memory_order_relaxed) : 0; - state_to_publish.media_lock_wait_count = d->media_lock_wait_count.load(std::memory_order_relaxed); - state_to_publish.media_lock_wait_total_ns = d->media_lock_wait_total_ns.load(std::memory_order_relaxed); - state_to_publish.media_lock_wait_max_ns = d->media_lock_wait_max_ns.load(std::memory_order_relaxed); - state_to_publish.readiness_check_count = d->readiness_check_count.load(std::memory_order_relaxed); - state_to_publish.readiness_reject_count = d->readiness_reject_count.load(std::memory_order_relaxed); - state_to_publish.close_join_total_ns = d->close_join_total_ns.load(std::memory_order_relaxed); - state_to_publish.close_join_max_ns = d->close_join_max_ns.load(std::memory_order_relaxed); - state_to_publish.current_close_join_ns = d->close_join_active.load(std::memory_order_acquire) - ? now - d->close_join_started_ns.load(std::memory_order_relaxed) : 0; - state_to_publish.track_open = d->callbacks->track_open.load(std::memory_order_relaxed); - state_to_publish.closed = d->callbacks->closed.load(std::memory_order_relaxed); - state_to_publish.sender_running = d->sender_loop_active.load(std::memory_order_relaxed); - d->state.advance(); - - const auto& published = *d->state.pending; + const auto queued_frame_count = + d->queued_frame_count.load(std::memory_order_relaxed); + const auto sent_frame_count = + d->sent_frame_count.load(std::memory_order_relaxed); + const auto send_total_ns = d->send_total_ns.load(std::memory_order_relaxed); + const auto readiness_check_count = + d->readiness_check_count.load(std::memory_order_relaxed); + const auto close_join_total_ns = + d->close_join_total_ns.load(std::memory_order_relaxed); const auto milliseconds = [](std::uint64_t nanoseconds) { return static_cast(nanoseconds) / 1'000'000.0; }; @@ -423,34 +351,30 @@ nlohmann::json WebRtc_Video_Session::diagnostics() const { {"kind", "webrtc_transport_state"}, {"protocol", "aethera.gallery.webrtc"}, {"version", 1}, - {"track_open", published.track_open}, - {"closed", published.closed}, - {"sender_running", published.sender_running}, - {"outstanding_video_frames", published.outstanding_video_frames}, - {"transport_buffered_bytes", published.transport_buffered_bytes}, - {"queued_frame_count", published.queued_frame_count}, - {"rejected_frame_count", published.rejected_frame_count}, - {"queued_megabytes", static_cast(published.queued_byte_count) / + {"track_open", d->callbacks->track_open.load(std::memory_order_relaxed)}, + {"closed", d->callbacks->closed.load(std::memory_order_relaxed)}, + {"sender_running", d->sender_loop_active.load(std::memory_order_relaxed)}, + {"outstanding_video_frames", outstanding_video_frames}, + {"transport_buffered_bytes", d->transport_buffered_bytes.load(std::memory_order_relaxed)}, + {"queued_frame_count", queued_frame_count}, + {"rejected_frame_count", d->rejected_frame_count.load(std::memory_order_relaxed)}, + {"queued_megabytes", static_cast(d->queued_byte_count.load(std::memory_order_relaxed)) / (1024.0 * 1024.0)}, - {"sent_frame_count", published.sent_frame_count}, - {"sent_megabytes", static_cast(published.sent_byte_count) / + {"sent_frame_count", sent_frame_count}, + {"sent_megabytes", static_cast(d->sent_byte_count.load(std::memory_order_relaxed)) / (1024.0 * 1024.0)}, - {"send_failure_count", published.send_failure_count}, - {"send_average_ms", published.sent_frame_count == 0 ? 0.0 : - milliseconds(published.send_total_ns) / - static_cast(published.sent_frame_count)}, - {"send_max_ms", milliseconds(published.send_max_ns)}, - {"current_send_ms", milliseconds(published.current_send_ns)}, - {"media_lock_wait_count", published.media_lock_wait_count}, - {"media_lock_wait_average_ms", published.media_lock_wait_count == 0 ? 0.0 : - milliseconds(published.media_lock_wait_total_ns) / - static_cast(published.media_lock_wait_count)}, - {"media_lock_wait_max_ms", milliseconds(published.media_lock_wait_max_ns)}, - {"readiness_check_count", published.readiness_check_count}, - {"readiness_reject_count", published.readiness_reject_count}, - {"close_join_total_ms", milliseconds(published.close_join_total_ns)}, - {"close_join_max_ms", milliseconds(published.close_join_max_ns)}, - {"current_close_join_ms", milliseconds(published.current_close_join_ns)}}; + {"send_failure_count", d->send_failure_count.load(std::memory_order_relaxed)}, + {"send_average_ms", sent_frame_count == 0 ? 0.0 : + milliseconds(send_total_ns) / static_cast(sent_frame_count)}, + {"send_max_ms", milliseconds(d->send_max_ns.load(std::memory_order_relaxed))}, + {"current_send_ms", milliseconds(d->send_active.load(std::memory_order_acquire) + ? now - d->send_started_ns.load(std::memory_order_relaxed) : 0)}, + {"readiness_check_count", readiness_check_count}, + {"readiness_reject_count", d->readiness_reject_count.load(std::memory_order_relaxed)}, + {"close_join_total_ms", milliseconds(close_join_total_ns)}, + {"close_join_max_ms", milliseconds(d->close_join_max_ns.load(std::memory_order_relaxed))}, + {"current_close_join_ms", milliseconds(d->close_join_active.load(std::memory_order_acquire) + ? now - d->close_join_started_ns.load(std::memory_order_relaxed) : 0)}}; } void WebRtc_Video_Session::close() noexcept { @@ -477,11 +401,11 @@ void WebRtc_Video_Session::close() noexcept { d->outstanding_video_frames.store(0, std::memory_order_release); std::lock_guard lock(d->media_mutex); - try { if (d->video_track) d->video_track->close(); } + const auto track = d->video_track.exchange({}, std::memory_order_acq_rel); + try { if (track) track->close(); } catch (...) {} try { if (d->peer) d->peer->close(); } catch (...) {} - d->video_track.reset(); d->rtp_config.reset(); d->peer.reset(); } diff --git a/web_server/src/Web_Server.cpp b/web_server/src/Web_Server.cpp index 80237c8..1b13d1b 100644 --- a/web_server/src/Web_Server.cpp +++ b/web_server/src/Web_Server.cpp @@ -3,7 +3,6 @@ #include "Gallery_WebSocket.hpp" #include "Graph_WebSocket.hpp" #include "Gallery_Plots.hpp" -#include "Page_Session.hpp" #include #include #include @@ -176,26 +175,15 @@ int run_web_server(std::uint16_t port, const std::filesystem::path& asset_root) }; add_media_groups("2d", std::move(gallery_2d)); add_media_groups("3d", std::move(gallery_3d)); - auto page_session = std::make_shared(); auto resolve_plot = [plots](std::string_view id) { return find_plot(*plots, id); }; - auto websocket = std::make_shared( - resolve_plot, page_session); + auto websocket = std::make_shared(resolve_plot); auto gallery_websocket = std::make_shared( [gallery_streams](std::string_view id) { const auto found = gallery_streams->find(std::string{id}); return found == gallery_streams->end() ? nullptr : found->second; - }, page_session); + }); auto& app = drogon::app(); - app.registerHandler("/page/session", [page_session]( - const drogon::HttpRequestPtr&, - std::function&& callback) { - callback(json_response({ - {"protocol", "aethera.page.session"}, - {"version", 1}, - {"token", page_session->begin_page()}})); - }, {drogon::Post}); - app.registerHandler("/plot", [plots, plot_media](const drogon::HttpRequestPtr&, std::function&& callback) { nlohmann::json result = nlohmann::json::array(); diff --git a/web_server/src/detail/Gallery_Frame_Atlas.cpp b/web_server/src/detail/Gallery_Frame_Atlas.cpp index a53e3ed..88ba2e3 100644 --- a/web_server/src/detail/Gallery_Frame_Atlas.cpp +++ b/web_server/src/detail/Gallery_Frame_Atlas.cpp @@ -3,19 +3,23 @@ #include #include #include -#include #include #include namespace aethera::web::detail { struct Gallery_Frame_Atlas::Private { struct Source { - mutable std::mutex mutex; /* 仅保护当前 Plot 的最近完成快照。 */ - std::shared_ptr latest_completion; /* 最近逻辑完成身份;像素可为空。 */ - std::shared_ptr latest_pixels; /* 最近可合成的实际 RGBA 画面。 */ - std::uint64_t composited_rendered_sequence{}; /* 图集像素当前包含的真实画面序号。 */ - std::uint64_t completion_count{}; /* 当前槽位接受的逻辑完成回调总数。 */ - std::uint64_t rendered_frame_count{}; /* 当前槽位接受的不同真实画面总数。 */ + struct Published { + std::shared_ptr latest_completion; + std::shared_ptr latest_pixels; + std::uint64_t completion_count{}; + std::uint64_t rendered_frame_count{}; + }; + /* Plot 完成线程发布不可变版本,图集线程只读取同一个版本;不再把四个 + * 相关字段分别复制出锁区。 */ + std::atomic> published{ + std::make_shared()}; + std::atomic_uint64_t composited_rendered_sequence{}; }; Gallery_Atlas_Description description{}; /* 尺寸与槽位布局的唯一权威描述。 */ @@ -83,7 +87,6 @@ Gallery_Frame_Atlas::Accept_Frame_Result Gallery_Frame_Atlas::accept_frame( return Accept_Frame_Result::invalid_frame; } auto& source = *d->sources[slot]; - std::lock_guard lock(source.mutex); if (frame->pixels) { const auto candidate = static_cast(frame->layout); int unknown{-1}; @@ -95,21 +98,29 @@ Gallery_Frame_Atlas::Accept_Frame_Result Gallery_Frame_Atlas::accept_frame( return Accept_Frame_Result::invalid_frame; } } - if (source.latest_completion && - frame->sequence <= source.latest_completion->sequence) { - d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed); - return Accept_Frame_Result::stale_frame; + auto current = source.published.load(std::memory_order_acquire); + for (;;) { + if (current->latest_completion && + frame->sequence <= current->latest_completion->sequence) { + d->rejected_frame_count.fetch_add(1, std::memory_order_relaxed); + return Accept_Frame_Result::stale_frame; + } + auto next = std::make_shared(*current); + if (frame->rendered_sequence != 0 && + (!current->latest_completion || + frame->rendered_sequence != + current->latest_completion->rendered_sequence)) + ++next->rendered_frame_count; + next->latest_completion = frame; + if (frame->pixels && frame->rendered_sequence != 0) + next->latest_pixels = frame; + ++next->completion_count; + std::shared_ptr desired = next; + if (source.published.compare_exchange_weak( + current, std::move(desired), std::memory_order_release, + std::memory_order_acquire)) + return Accept_Frame_Result::accepted; } - if (frame->rendered_sequence != 0 && - (!source.latest_completion || - frame->rendered_sequence != - source.latest_completion->rendered_sequence)) - ++source.rendered_frame_count; - source.latest_completion = frame; - if (frame->pixels && frame->rendered_sequence != 0) - source.latest_pixels = std::move(frame); - ++source.completion_count; - return Accept_Frame_Result::accepted; } Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { @@ -126,17 +137,9 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { static_cast(d->description.tile_width) * 4U; for (std::size_t slot = 0; slot < d->sources.size(); ++slot) { auto& source = *d->sources[slot]; - std::shared_ptr latest_completion; - std::shared_ptr latest_pixels; - std::uint64_t completion_count{}; - std::uint64_t rendered_frame_count{}; - { - std::lock_guard lock(source.mutex); - latest_completion = source.latest_completion; - latest_pixels = source.latest_pixels; - completion_count = source.completion_count; - rendered_frame_count = source.rendered_frame_count; - } + const auto published = source.published.load(std::memory_order_acquire); + const auto& latest_completion = published->latest_completion; + const auto& latest_pixels = published->latest_pixels; if (!latest_completion) { ++result.missing_tile_count; result.sources.push_back({}); @@ -146,7 +149,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { ++result.missing_tile_count; } else if (latest_pixels->rendered_sequence != - source.composited_rendered_sequence) { + source.composited_rendered_sequence.load( + std::memory_order_acquire)) { const auto& layout = d->description.sources[slot]; for (std::uint32_t y = 0; y < d->description.tile_height; ++y) { const auto source_offset = @@ -160,8 +164,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { latest_pixels->pixels->data() + source_offset, source_row_bytes); } - source.composited_rendered_sequence = - latest_pixels->rendered_sequence; + source.composited_rendered_sequence.store( + latest_pixels->rendered_sequence, std::memory_order_release); ++result.fresh_tile_count; } result.sources.push_back({ @@ -169,8 +173,8 @@ Gallery_Atlas_Composition Gallery_Frame_Atlas::compose() { latest_completion->correlation_id, latest_completion->rendered_sequence, latest_completion->rendered_correlation_id, - completion_count, - rendered_frame_count}); + published->completion_count, + published->rendered_frame_count}); } result.pixels = d->pixels; return result; diff --git a/web_server/src/detail/Gallery_Frame_Clock.cpp b/web_server/src/detail/Gallery_Frame_Clock.cpp index 9ccab16..5a52d0f 100644 --- a/web_server/src/detail/Gallery_Frame_Clock.cpp +++ b/web_server/src/detail/Gallery_Frame_Clock.cpp @@ -13,10 +13,13 @@ struct Gallery_Frame_Clock::Private { std::chrono::nanoseconds interval{}; Tick_Handler tick_handler{}; Failure_Handler failure_handler{}; + /* condition_variable_any 的 stop_token 等待和 jthread 替换必须共享生命周期 + * 临界区;它不是业务状态锁,双缓冲不能替代等待协议。 */ std::mutex lifecycle_mutex{}; std::condition_variable_any wake{}; std::jthread thread{}; std::atomic_bool running{}; + std::atomic_uint64_t generation{}; /* 区分 stop 期间接管的新时钟。 */ Private(double frame_rate, Tick_Handler value_tick_handler, Failure_Handler value_failure_handler) @@ -32,14 +35,16 @@ struct Gallery_Frame_Clock::Private { throw std::invalid_argument("gallery frame clock requires a tick handler"); } - void report(std::exception_ptr failure) noexcept { - running.store(false, std::memory_order_release); + void report(std::uint64_t active_generation, + std::exception_ptr failure) noexcept { + if (generation.load(std::memory_order_acquire) == active_generation) + running.store(false, std::memory_order_release); if (!failure_handler) return; try { failure_handler(std::move(failure)); } catch (...) {} } - void run(std::stop_token stop) noexcept { + void run(std::stop_token stop, std::uint64_t active_generation) noexcept { const auto origin = std::chrono::steady_clock::now(); std::uint64_t sequence{}; auto deadline = origin + interval; @@ -60,11 +65,12 @@ struct Gallery_Frame_Clock::Private { deadline += interval * ((now - deadline) / interval + 1); } catch (...) { - report(std::current_exception()); + report(active_generation, std::current_exception()); return; } } - running.store(false, std::memory_order_release); + if (generation.load(std::memory_order_acquire) == active_generation) + running.store(false, std::memory_order_release); } }; @@ -79,17 +85,25 @@ Gallery_Frame_Clock::~Gallery_Frame_Clock() { stop(); } void Gallery_Frame_Clock::start() { std::lock_guard lock(d->lifecycle_mutex); if (d->running.exchange(true, std::memory_order_acq_rel)) return; - d->thread = std::jthread([state = d](std::stop_token stop) { - state->run(stop); + const auto generation = d->generation.fetch_add( + 1, std::memory_order_acq_rel) + 1; + d->thread = std::jthread([state = d, generation](std::stop_token stop) { + state->run(stop, generation); }); } void Gallery_Frame_Clock::stop() noexcept { - if (!d || !d->running.exchange(false, std::memory_order_acq_rel)) return; - d->thread.request_stop(); - d->wake.notify_all(); - if (d->thread.joinable() && - d->thread.get_id() != std::this_thread::get_id()) - d->thread.join(); + if (!d) return; + std::jthread thread; + { + std::lock_guard lock(d->lifecycle_mutex); + if (!d->running.exchange(false, std::memory_order_acq_rel)) return; + d->thread.request_stop(); + d->wake.notify_all(); + thread = std::move(d->thread); + } + if (!thread.joinable()) return; + if (thread.get_id() == std::this_thread::get_id()) thread.detach(); + else thread.join(); } } diff --git a/web_server/tests/Page_Session_Tests.cpp b/web_server/tests/Page_Session_Tests.cpp deleted file mode 100644 index d32eb91..0000000 --- a/web_server/tests/Page_Session_Tests.cpp +++ /dev/null @@ -1,36 +0,0 @@ -#include "web_server/src/Page_Session.hpp" -#include - -namespace aethera::web { -TEST(Page_Session, New_Page_Displaces_Previous_Page_Without_Waiting_For_Close) { - Page_Session session; - const auto first_token = session.begin_page(); - std::size_t displaced{}; - const auto first_connection = session.attach( - first_token, [&displaced] { ++displaced; }); - ASSERT_TRUE(first_connection.has_value()); - - const auto second_token = session.begin_page(); - EXPECT_NE(first_token, second_token); - EXPECT_EQ(displaced, 1U); - EXPECT_EQ(session.attach(first_token, [] {}).error(), - Page_Session::Attach_Result::stale_session); - - const auto second_connection = session.attach(second_token, [] {}); - EXPECT_TRUE(second_connection.has_value()); - session.detach(*first_connection); - session.detach(*second_connection); -} - -TEST(Page_Session, One_Page_Can_Attach_All_Of_Its_Connections) { - Page_Session session; - const auto token = session.begin_page(); - const auto gallery = session.attach(token, [] {}); - const auto plot = session.attach(token, [] {}); - ASSERT_TRUE(gallery.has_value()); - ASSERT_TRUE(plot.has_value()); - EXPECT_NE(*gallery, *plot); - session.detach(*gallery); - session.detach(*plot); -} -} diff --git a/web_server/tests/Sliding_Statistics_Tests.cpp b/web_server/tests/Sliding_Statistics_Tests.cpp index 5baff2d..bbb804a 100644 --- a/web_server/tests/Sliding_Statistics_Tests.cpp +++ b/web_server/tests/Sliding_Statistics_Tests.cpp @@ -1,5 +1,6 @@ #include #include +#include #include namespace aethera::web { @@ -52,4 +53,18 @@ TEST(Sliding_Statistics, Estimates_Quantiles_Without_Growing_The_Window) { EXPECT_NEAR(state.p95, 950.0, 15.0); EXPECT_NEAR(state.p99, 990.0, 15.0); } + +TEST(Sliding_Statistics, Keeps_Approximate_Quantiles_Ordered_For_Nonstationary_Input) { + Sliding_Statistics statistics{16}; + for (std::size_t phase = 0; phase < 200; ++phase) { + const double baseline = phase % 2 == 0 ? 1'000'000.0 : -1'000'000.0; + for (std::size_t sample = 0; sample < 17; ++sample) { + const auto state = statistics.submit( + baseline + static_cast(sample * sample)); + EXPECT_TRUE(std::isfinite(state.trimmed_average)); + EXPECT_LE(state.p50, state.p95); + EXPECT_LE(state.p95, state.p99); + } + } +} } diff --git a/webapp_gallery/src/app.tsx b/webapp_gallery/src/app.tsx index 8296321..57943bd 100644 --- a/webapp_gallery/src/app.tsx +++ b/webapp_gallery/src/app.tsx @@ -1,4 +1,4 @@ -import {createContext, memo, useCallback, useContext, useEffect, useMemo, useRef, useState} from "react"; +import {memo, useCallback, useEffect, useMemo, useRef, useState} from "react"; import {I18nLabel, Layout, Model, type IJsonModel, type TabNode} from "flexlayout-react"; import {Responsive, useContainerWidth, type LayoutItem, type ResponsiveLayouts} from "react-grid-layout"; import {ResizableBox} from "react-resizable"; @@ -102,12 +102,10 @@ type Taskflow_Runtime_State = {protocol: "aethera.taskflow.runtime"; version: 1; task_types: Taskflow_Type_State[]; workers: Taskflow_Worker_State[]}; const default_plot_execution_policy = (): Plot_Execution_Policy => ({visible: true}); -const Page_Session_Context = createContext(""); -function socket_url(path: string, page_session: string) { +function socket_url(path: string) { const url = new URL(path, location.href); url.protocol = location.protocol === "https:" ? "wss:" : "ws:"; - url.searchParams.set("session", page_session); return url.toString(); } @@ -194,7 +192,7 @@ const connecting_gallery_video = (): Gallery_Video_State => ({ playback: empty_video_playback(), transport: null, error: null }); -function use_gallery_videos(plots: Plot[], page_session: string): Gallery_Video_States { +function use_gallery_videos(plots: Plot[]): 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({}); @@ -227,7 +225,7 @@ function use_gallery_videos(plots: Plot[], page_session: string): Gallery_Video_ }; window.addEventListener("aethera:gallery-diagnostics", receive_server_diagnostics); - const socket = new ReconnectingWebSocket(socket_url(endpoint, page_session), [], { + const socket = new ReconnectingWebSocket(socket_url(endpoint), [], { minReconnectionDelay: 300, maxReconnectionDelay: 5000, reconnectionDelayGrowFactor: 1.6, maxRetries: Number.POSITIVE_INFINITY }); @@ -311,16 +309,9 @@ function use_gallery_videos(plots: Plot[], page_session: string): Gallery_Video_ create_peer(); update(current => ({...current, status: "CONNECTING", error: null})); }; - socket.onclose = event => { + socket.onclose = () => { peer?.close(); peer = null; - if (event.code === 1008) { - stopped = true; - update(current => ({...current, status: "OFFLINE", stream: null, - error: "该页面会话已失效"})); - socket.close(4002, "Page session rejected"); - return; - } if (!stopped) update(current => ({...current, status: "CONNECTING", stream: null})); }; socket.onmessage = event => { @@ -329,16 +320,6 @@ function use_gallery_videos(plots: Plot[], page_session: string): Gallery_Video_ try { decoded = JSON.parse(event.data); } catch { return; } if (!decoded || typeof decoded !== "object") return; const message = decoded as Record; - if (message.kind === "page_session_replaced" || - message.kind === "page_session_rejected") { - stopped = true; - peer?.close(); - peer = null; - update(current => ({...current, status: "OFFLINE", stream: null, - error: "该页面会话已被更新的页面接管"})); - socket.close(4002, "Page session replaced"); - return; - } if (valid_gallery_layout(decoded)) { update(current => ({...current, layout: decoded})); window.dispatchEvent(new CustomEvent("aethera:gallery-layout", @@ -379,7 +360,7 @@ function use_gallery_videos(plots: Plot[], page_session: string): Gallery_Video_ }; }); return () => cleanups.forEach(cleanup => cleanup()); - }, [media_key, page_session]); + }, [media_key]); return states; } @@ -387,7 +368,6 @@ function use_plot_stream(plot: Plot, surface_ref: React.RefObject, shared_playback: Video_Playback_Metrics, tile_width = 720, tile_height = 420) { - const page_session = useContext(Page_Session_Context); const [status, set_status] = useState("CONNECTING"); const [metrics, set_metrics] = useState(null); const [server_diagnostics, set_server_diagnostics] = useState(null); @@ -447,8 +427,7 @@ function use_plot_stream(plot: Plot, useEffect(() => { let stopped = false; let samples: Frame_Sample[] = []; - const socket = new ReconnectingWebSocket( - socket_url(plot.websocket, page_session), [], { + const socket = new ReconnectingWebSocket(socket_url(plot.websocket), [], { minReconnectionDelay: 300, maxReconnectionDelay: 5000, reconnectionDelayGrowFactor: 1.6, maxRetries: Number.POSITIVE_INFINITY }); @@ -531,29 +510,12 @@ function use_plot_stream(plot: Plot, set_error(null); socket.send(JSON.stringify({kind: "stream", viewport: viewport_ref.current})); }; - socket.onclose = event => { - if (event.code === 1008) { - stopped = true; - set_status("OFFLINE"); - set_error("该页面会话已失效"); - socket.close(4002, "Page session rejected"); - return; - } - if (!stopped) set_status("CONNECTING"); - }; + socket.onclose = () => { if (!stopped) set_status("CONNECTING"); }; socket.onmessage = event => { if (typeof event.data !== "string") return; let decoded: unknown; try { decoded = JSON.parse(event.data); } catch { return; } const message = decoded as {kind?: string; message?: string}; - if (message?.kind === "page_session_replaced" || - message?.kind === "page_session_rejected") { - stopped = true; - set_status("OFFLINE"); - set_error("该页面会话已被更新的页面接管"); - socket.close(4002, "Page session replaced"); - return; - } if (message?.kind === "plot_error") { stopped = true; set_status("OFFLINE"); @@ -571,7 +533,7 @@ function use_plot_stream(plot: Plot, socket_ref.current = null; socket.close(); }; - }, [plot.id, plot.websocket, plot.diagnostics, page_session, surface_ref, transmit, + }, [plot.id, plot.websocket, plot.diagnostics, surface_ref, transmit, read_browser_input_statistics]); useEffect(() => { @@ -1813,10 +1775,10 @@ function load_workspace_model() { catch { localStorage.removeItem(workspace_layout_key); return Model.fromJson(default_workspace_layout); } } -export function App({page_session}: {page_session: string}) { +export function App() { const [plots, set_plots] = useState([]); const [category, set_category] = useState("全部"); const [selected, set_selected] = useState(null); use_selected_plot_diagnostics(selected); - const gallery_videos = use_gallery_videos(plots, page_session); + const gallery_videos = use_gallery_videos(plots); const [execution_policies, set_execution_policies] = useState({}); const [schema, set_schema] = useState(null); const [state_histories, set_state_histories] = useState({}); @@ -1913,8 +1875,6 @@ export function App({page_session}: {page_session: string}) { if (node.getComponent() === "taskflow-frame") return ; return
未知工作区面板。
; }; - return -
layout_labels[label]} - onModelChange={model => localStorage.setItem(workspace_layout_key, JSON.stringify(model.toJson()))}/>
-
; + return
layout_labels[label]} + onModelChange={model => localStorage.setItem(workspace_layout_key, JSON.stringify(model.toJson()))}/>
; } diff --git a/webapp_gallery/src/main.tsx b/webapp_gallery/src/main.tsx index 7b7581b..b099816 100644 --- a/webapp_gallery/src/main.tsx +++ b/webapp_gallery/src/main.tsx @@ -2,27 +2,4 @@ import {StrictMode} from "react"; import {createRoot} from "react-dom/client"; import {App} from "./app"; import "./styles.css"; - -type Page_Session_Response = { - protocol: "aethera.page.session"; - version: 1; - token: string; -}; - -async function begin_page() { - const response = await fetch("/page/session", {method: "POST"}); - if (!response.ok) throw new Error(`Page session failed: HTTP ${response.status}`); - const session = await response.json() as Partial; - if (session.protocol !== "aethera.page.session" || session.version !== 1 || - typeof session.token !== "string" || session.token.length === 0) - throw new Error("Page session response is invalid"); - return session.token; -} - -const root = createRoot(document.getElementById("root")!); -void begin_page().then(page_session => { - root.render(); -}).catch(failure => { - const message = failure instanceof Error ? failure.message : "Page session failed"; - root.render(
{message}
); -}); +createRoot(document.getElementById("root")!).render();