From 513f838874d74211b845ef0760fb7d99c048d257 Mon Sep 17 00:00:00 2001 From: wyc <1104749580@qq.com> Date: Wed, 26 Aug 2026 11:29:46 +0800 Subject: [PATCH] =?UTF-8?q?=E5=8A=A0=E4=B8=8A=E5=85=A8=E5=B1=80=E8=B0=83?= =?UTF-8?q?=E5=BA=A6=E5=99=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- kernel/src/kernel/Frame_Pacing_Policy.cpp | 89 +++++ kernel/src/kernel/Frame_Pacing_Policy.hpp | 57 +++ kernel/src/kernel/Frame_Scheduler.cpp | 328 ++++++++++++++++++ kernel/src/kernel/Frame_Scheduler.hpp | 69 ++++ .../detail/Gpu_Completion_Service.cpp | 298 ++++++---------- .../detail/Gpu_Completion_Service.hpp | 2 +- .../detail/Gpu_Completion_Service.ipp | 25 +- web_server/src/Gallery_Video_Stream.cpp | 36 +- web_server/src/Gallery_Video_Stream.ipp | 3 - web_server/src/Plot.cpp | 147 ++++---- web_server/src/Plot.hpp | 2 +- web_server/src/detail/Gallery_Frame_Clock.cpp | 109 ------ web_server/src/detail/Gallery_Frame_Clock.hpp | 22 -- 13 files changed, 737 insertions(+), 450 deletions(-) create mode 100644 kernel/src/kernel/Frame_Pacing_Policy.cpp create mode 100644 kernel/src/kernel/Frame_Pacing_Policy.hpp create mode 100644 kernel/src/kernel/Frame_Scheduler.cpp create mode 100644 kernel/src/kernel/Frame_Scheduler.hpp delete mode 100644 web_server/src/detail/Gallery_Frame_Clock.cpp delete mode 100644 web_server/src/detail/Gallery_Frame_Clock.hpp diff --git a/kernel/src/kernel/Frame_Pacing_Policy.cpp b/kernel/src/kernel/Frame_Pacing_Policy.cpp new file mode 100644 index 0000000..8688e8e --- /dev/null +++ b/kernel/src/kernel/Frame_Pacing_Policy.cpp @@ -0,0 +1,89 @@ +#include "Frame_Pacing_Policy.hpp" +#include +#include +#include + +namespace aethera { + +Frame_Pacing_Policy::Frame_Pacing_Policy() + : properties_(std::make_shared()) {} + +Frame_Pacing_Properties Frame_Pacing_Policy::read() const { + return *properties_.load(std::memory_order_acquire); +} + +template +void Frame_Pacing_Policy::update(Edit&& edit) { + auto current = properties_.load(std::memory_order_acquire); + for (;;) { + auto next = std::make_shared(*current); + std::forward(edit)(*next); + std::shared_ptr desired = next; + if (properties_.compare_exchange_weak( + current, desired, std::memory_order_release, + std::memory_order_acquire)) + return; + } +} + +void Frame_Pacing_Policy::set_render_enabled(bool enabled) { + update([enabled](auto& value) { value.render_enabled = enabled; }); +} + +void Frame_Pacing_Policy::set_video_enabled(bool enabled) { + update([enabled](auto& value) { value.video_enabled = enabled; }); +} + +void Frame_Pacing_Policy::set_mode(Frame_Pacing_Mode mode) { + update([mode](auto& value) { value.mode = mode; }); +} + +void Frame_Pacing_Policy::set_fixed_rate(double fps) { + if (!std::isfinite(fps) || fps <= 0.0) + throw std::invalid_argument("frame pacing fps must be positive"); + update([fps](auto& value) { value.fixed_rate_fps = fps; }); +} + +double Frame_Pacing_Policy::scheduled_rate_fps() const { + const auto value = read(); + if (!value.render_enabled || value.mode == Frame_Pacing_Mode::manual) + return 0.0; + return value.mode == Frame_Pacing_Mode::maximum_rate + ? value.maximum_rate_fps : value.fixed_rate_fps; +} + +bool Frame_Pacing_Policy::accept_periodic_tick(double time_milliseconds) { + static_cast(time_milliseconds); + const auto value = read(); + if (!value.render_enabled || value.mode == Frame_Pacing_Mode::manual) + return false; + if (skip_next_periodic_.exchange(false, std::memory_order_acq_rel)) + return false; + return !in_flight_.load(std::memory_order_acquire); +} + +bool Frame_Pacing_Policy::request_immediate() { + const auto value = read(); + if (!value.render_enabled) return false; + skip_next_periodic_.store(true, std::memory_order_release); + if (in_flight_.load(std::memory_order_acquire)) { + urgent_pending_.store(true, std::memory_order_release); + return false; + } + return true; +} + +void Frame_Pacing_Policy::frame_submitted() { + in_flight_.store(true, std::memory_order_release); +} + +bool Frame_Pacing_Policy::frame_completed() { + in_flight_.store(false, std::memory_order_release); + return urgent_pending_.exchange(false, std::memory_order_acq_rel); +} + +void Frame_Pacing_Policy::frame_rejected() { + in_flight_.store(false, std::memory_order_release); +} + +} diff --git a/kernel/src/kernel/Frame_Pacing_Policy.hpp b/kernel/src/kernel/Frame_Pacing_Policy.hpp new file mode 100644 index 0000000..2c30f28 --- /dev/null +++ b/kernel/src/kernel/Frame_Pacing_Policy.hpp @@ -0,0 +1,57 @@ +#pragma once +#include +#include +#include + +namespace aethera { + +enum struct Frame_Pacing_Mode : std::uint8_t { + manual, + fixed_rate, + maximum_rate +}; + +struct Frame_Pacing_Properties { + bool render_enabled{true}; + bool video_enabled{true}; + Frame_Pacing_Mode mode{Frame_Pacing_Mode::fixed_rate}; + double fixed_rate_fps{100.0}; + double maximum_rate_fps{100.0}; +}; + +/* + * 每个 Scene/Plot 独立持有的帧策略。 + * Scheduler 只负责 WHEN;这里负责 WHETHER。 + */ +struct Frame_Pacing_Policy final { + Frame_Pacing_Policy(); + + [[nodiscard]] Frame_Pacing_Properties read() const; + void set_render_enabled(bool enabled); + void set_video_enabled(bool enabled); + void set_mode(Frame_Pacing_Mode mode); + void set_fixed_rate(double fps); + + /* 当前策略应向全局 Frame_Scheduler 注册的频率;0 表示事件驱动/关闭周期。 */ + [[nodiscard]] double scheduled_rate_fps() const; + + /* + * 周期 tick 是否允许生成新帧。immediate 会提前生成一帧并消费下一次周期机会。 + * in_flight 只表达 Scene 是否已有一帧尚未完成,不改变 Scene::render(Frame*) 抽象。 + */ + [[nodiscard]] bool accept_periodic_tick(double time_milliseconds); + [[nodiscard]] bool request_immediate(); + void frame_submitted(); + [[nodiscard]] bool frame_completed(); + void frame_rejected(); + +private: + template void update(Edit&& edit); + + std::atomic> properties_; + std::atomic_bool in_flight_{}; + std::atomic_bool skip_next_periodic_{}; + std::atomic_bool urgent_pending_{}; +}; + +} diff --git a/kernel/src/kernel/Frame_Scheduler.cpp b/kernel/src/kernel/Frame_Scheduler.cpp new file mode 100644 index 0000000..8dd71eb --- /dev/null +++ b/kernel/src/kernel/Frame_Scheduler.cpp @@ -0,0 +1,328 @@ +#include "Frame_Scheduler.hpp" +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace aethera { +namespace { +using Scheduler_Clock = Frame_Scheduler::Clock; +using Nanoseconds = std::chrono::nanoseconds; + +::Tick duration_to_ticks(Nanoseconds duration) noexcept { + if (duration <= Nanoseconds::zero()) return 1; + const auto quantum = std::chrono::duration_cast( + Frame_Scheduler::tick_duration).count(); + const auto value = duration.count(); + const auto ticks = (value + quantum - 1) / quantum; + return static_cast<::Tick>(std::max(1, ticks)); +} + +Nanoseconds fps_period(double fps) { + if (!std::isfinite(fps) || fps <= 0.0) + throw std::invalid_argument("frame scheduler fps must be positive"); + const auto count = static_cast( + std::llround(1'000'000'000.0 / fps)); + if (count <= 0) + throw std::invalid_argument("frame scheduler fps exceeds nanosecond resolution"); + return Nanoseconds{count}; +} +} + +struct Frame_Scheduler::Timer::State { + enum struct Mode : std::uint8_t { stopped, once, periodic }; + + Frame_Scheduler::Private* owner{}; + std::uint64_t id{}; + Handler handler{}; + std::atomic_bool alive{true}; + Mode mode{Mode::stopped}; /* 仅 Scheduler thread 访问。 */ + Nanoseconds period{}; /* 仅 Scheduler thread 访问。 */ + Scheduler_Clock::time_point next_deadline{}; /* 真实 deadline,不按 tick 累积。 */ + TimerEvent> event; + + State(Frame_Scheduler::Private* value_owner, std::uint64_t value_id, + Handler value_handler); +}; + +struct Frame_Scheduler::Private { + enum struct Command_Type : std::uint8_t { + periodic, + once, + cancel, + destroy + }; + struct Command { + Command_Type type{}; + std::shared_ptr state{}; + Nanoseconds value{}; + std::shared_ptr> completion{}; + }; + + TimerWheel wheel{}; + Scheduler_Clock::time_point origin{Scheduler_Clock::now()}; + Scheduler_Clock::time_point last_advance{origin}; + std::mutex mutex{}; + std::condition_variable wake{}; + std::deque commands{}; + std::unordered_map> timers{}; + std::atomic_uint64_t next_id{1}; + std::atomic_bool stopping{}; + std::thread::id control_thread_id{}; + std::jthread thread{}; + + Private() : thread([this](std::stop_token stop) { run(stop); }) {} + + ~Private() { + stopping.store(true, std::memory_order_release); + thread.request_stop(); + wake.notify_all(); + } + + void enqueue(Command command) { + { + std::lock_guard lock(mutex); + commands.push_back(std::move(command)); + } + wake.notify_one(); + } + + std::shared_ptr make_timer(Handler handler) { + if (!handler) throw std::invalid_argument("frame scheduler timer handler is empty"); + const auto id = next_id.fetch_add(1, std::memory_order_relaxed); + auto state = std::make_shared(this, id, std::move(handler)); + { + std::lock_guard lock(mutex); + /* 注册表只用于维持 TimerEvent 到异步 cancel/destroy 被处理为止的生命周期。 */ + timers.emplace(id, state); + } + return state; + } + + + ::Tick delay_ticks(Scheduler_Clock::time_point deadline, + Scheduler_Clock::time_point now) const noexcept { + return duration_to_ticks(std::chrono::duration_cast(deadline - now)); + } + + void schedule_state(Timer::State& state, + Scheduler_Clock::time_point now) { + if (!state.alive.load(std::memory_order_acquire) || + state.mode == Timer::State::Mode::stopped) + return; + state.event.cancel(); + wheel.schedule(&state.event, delay_ticks(state.next_deadline, now)); + } + + void process_commands(Scheduler_Clock::time_point now) { + std::deque local; + { + std::lock_guard lock(mutex); + local.swap(commands); + } + for (auto& command : local) { + auto& state = command.state; + if (!state) continue; + switch (command.type) { + case Command_Type::periodic: + if (!state->alive.load(std::memory_order_acquire)) break; + state->mode = Timer::State::Mode::periodic; + state->period = command.value; + state->next_deadline = now + state->period; + schedule_state(*state, now); + break; + case Command_Type::once: + if (!state->alive.load(std::memory_order_acquire)) break; + state->mode = Timer::State::Mode::once; + state->period = {}; + state->next_deadline = now + std::max(command.value, Nanoseconds{1}); + schedule_state(*state, now); + break; + case Command_Type::cancel: + state->event.cancel(); + state->mode = Timer::State::Mode::stopped; + break; + case Command_Type::destroy: + state->event.cancel(); + state->mode = Timer::State::Mode::stopped; + { + std::lock_guard lock(mutex); + timers.erase(state->id); + } + break; + } + if (command.completion) { + try { command.completion->set_value(); } + catch (...) {} + } + } + } + + void fire(Timer::State& state) noexcept { + if (!state.alive.load(std::memory_order_acquire)) return; + const auto now = Scheduler_Clock::now(); + const ::Tick sequence = wheel.now(); + try { + state.handler(Frame_Scheduler::Tick{ + now, + static_cast(sequence), + std::chrono::duration(now - origin).count()}); + } + catch (...) { + /* Timer callback 是控制面入口,不能因为业务异常终止全局时钟。 */ + } + if (!state.alive.load(std::memory_order_acquire)) return; + if (state.mode == Timer::State::Mode::once) { + state.mode = Timer::State::Mode::stopped; + return; + } + if (state.mode != Timer::State::Mode::periodic || + state.period <= Nanoseconds::zero()) + return; + + state.next_deadline += state.period; + if (state.next_deadline <= now) { + const auto missed = (now - state.next_deadline) / state.period + 1; + state.next_deadline += state.period * missed; + } + schedule_state(state, now); + } + + void advance_to(Scheduler_Clock::time_point now) { + const auto elapsed = std::chrono::duration_cast(now - last_advance); + const auto quantum = std::chrono::duration_cast( + Frame_Scheduler::tick_duration); + if (elapsed < quantum) return; + const auto ticks = static_cast<::Tick>(elapsed / quantum); + bool finished = wheel.advance(ticks, 256); + while (!finished) finished = wheel.advance(0, 256); + last_advance += quantum * static_cast(ticks); + } + + void run(std::stop_token stop) noexcept { + control_thread_id = std::this_thread::get_id(); + try { + while (!stop.stop_requested() && + !stopping.load(std::memory_order_acquire)) { + auto now = Scheduler_Clock::now(); + process_commands(now); + advance_to(now); + + std::unique_lock lock(mutex); + if (stop.stop_requested() || + stopping.load(std::memory_order_acquire)) + break; + if (!commands.empty()) continue; + + /* timers.size() 包含已取消但尚未 destroy 的句柄;用 wheel 查询决定实际等待。 */ + constexpr ::Tick max_sleep_ticks = 10'000; /* 最多 1 s,便于可靠处理停止。 */ + const ::Tick next = wheel.ticks_to_next_event(max_sleep_ticks); + const auto deadline = last_advance + + Frame_Scheduler::tick_duration * static_cast(next); + wake.wait_until(lock, deadline, [this, &stop] { + return !commands.empty() || stop.stop_requested() || + stopping.load(std::memory_order_acquire); + }); + } + } + catch (...) { + stopping.store(true, std::memory_order_release); + } + } +}; + +Frame_Scheduler::Timer::State::State( + Frame_Scheduler::Private* value_owner, std::uint64_t value_id, + Handler value_handler) + : owner(value_owner), id(value_id), handler(std::move(value_handler)), + event([this] { + if (owner) owner->fire(*this); + }) {} + +Frame_Scheduler::Frame_Scheduler() : d(std::make_unique()) {} +Frame_Scheduler::~Frame_Scheduler() = default; + +Frame_Scheduler& Frame_Scheduler::instance() { + static Frame_Scheduler scheduler; + return scheduler; +} + +Frame_Scheduler::Timer Frame_Scheduler::make_timer(Handler handler) { + return Timer(d->make_timer(std::move(handler))); +} + +Frame_Scheduler::Timer::Timer(std::shared_ptr state) noexcept + : state_(std::move(state)) {} + +Frame_Scheduler::Timer::~Timer() noexcept { + if (!state_) return; + state_->alive.store(false, std::memory_order_release); + if (state_->owner) + state_->owner->enqueue({Private::Command_Type::destroy, state_, {}}); +} + +Frame_Scheduler::Timer::Timer(Timer&& other) noexcept + : state_(std::exchange(other.state_, {})) {} + +Frame_Scheduler::Timer& Frame_Scheduler::Timer::operator=(Timer&& other) noexcept { + if (this == &other) return *this; + this->~Timer(); + new (this) Timer(std::move(other)); + return *this; +} + +void Frame_Scheduler::Timer::start_periodic(double fps) { + if (!state_) throw std::logic_error("frame scheduler timer is empty"); + if (fps <= 0.0) { + cancel(); + return; + } + state_->alive.store(true, std::memory_order_release); + state_->owner->enqueue({Private::Command_Type::periodic, state_, fps_period(fps)}); +} + +void Frame_Scheduler::Timer::start_once(std::chrono::nanoseconds delay) { + if (!state_) throw std::logic_error("frame scheduler timer is empty"); + state_->alive.store(true, std::memory_order_release); + state_->owner->enqueue({Private::Command_Type::once, state_, delay}); +} + +void Frame_Scheduler::Timer::cancel() noexcept { + if (!state_ || !state_->owner) return; + state_->alive.store(false, std::memory_order_release); + state_->owner->enqueue({Private::Command_Type::cancel, state_, {}}); +} + + +void Frame_Scheduler::Timer::cancel_and_wait() noexcept { + if (!state_ || !state_->owner) return; + state_->alive.store(false, std::memory_order_release); + auto* owner = state_->owner; + if (owner->control_thread_id == std::this_thread::get_id()) { + state_->event.cancel(); + state_->mode = State::Mode::stopped; + return; + } + try { + auto completion = std::make_shared>(); + auto future = completion->get_future(); + owner->enqueue({Private::Command_Type::cancel, state_, {}, completion}); + future.wait(); + } + catch (...) {} +} + +bool Frame_Scheduler::Timer::valid() const noexcept { + return static_cast(state_); +} + +} diff --git a/kernel/src/kernel/Frame_Scheduler.hpp b/kernel/src/kernel/Frame_Scheduler.hpp new file mode 100644 index 0000000..27d5054 --- /dev/null +++ b/kernel/src/kernel/Frame_Scheduler.hpp @@ -0,0 +1,69 @@ +#pragma once +#include +#include +#include +#include + +namespace aethera { + +/* + * 全局轻量时间调度器。 + * + * - 整个进程只有一个控制线程。 + * - 底层使用 Ratas hierarchical timer wheel。 + * - callback 在控制线程执行,因此 callback 必须只做轻量控制工作, + * 真正 CPU 工作应继续投递到 Taskflow。 + * - Timer 是 RAII 句柄;每个 Scene/服务只持有轻量 Timer,不持有线程。 + */ +struct Frame_Scheduler final { + using Clock = std::chrono::steady_clock; + + struct Tick { + Clock::time_point issued_at{}; + std::uint64_t sequence{}; /* 全局时间轮 tick,所有 2D/3D Scene 共用同一时间轴。 */ + double time_milliseconds{}; /* 从全局 Scheduler origin 开始的单调毫秒。 */ + }; + + using Handler = std::function; + + struct Timer final { + public: + Timer() = default; + ~Timer() noexcept; + Timer(const Timer&) = delete; + Timer& operator=(const Timer&) = delete; + Timer(Timer&& other) noexcept; + Timer& operator=(Timer&& other) noexcept; + + /* 固定频率;fps <= 0 等价于 cancel。 */ + void start_periodic(double fps); + /* 单次延迟;delay <= 0 会在下一个 Scheduler 周期尽快触发。 */ + void start_once(std::chrono::nanoseconds delay); + void cancel() noexcept; + void cancel_and_wait() noexcept; /* 销毁拥有 callback 的对象前使用,保证控制线程不再执行该 callback。 */ + [[nodiscard]] bool valid() const noexcept; + + private: + struct State; + explicit Timer(std::shared_ptr state) noexcept; + std::shared_ptr state_{}; + friend struct Frame_Scheduler; + }; + + [[nodiscard]] static Frame_Scheduler& instance(); + [[nodiscard]] Timer make_timer(Handler handler); + + /* Ratas 量化粒度;策略层仍保存真实 deadline,避免周期累计漂移。 */ + static constexpr std::chrono::microseconds tick_duration{100}; + +private: + Frame_Scheduler(); + ~Frame_Scheduler(); + Frame_Scheduler(const Frame_Scheduler&) = delete; + Frame_Scheduler& operator=(const Frame_Scheduler&) = delete; + + struct Private; + std::unique_ptr d; +}; + +} diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp index 30a36fb..043fafc 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.cpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.cpp @@ -24,12 +24,7 @@ Gpu_Completion_Service::Private::Private() = default; Gpu_Completion_Service::Private::~Private() { stopping.store(true, std::memory_order_release); - { - std::lock_guard lock(service_mutex); - ++wake_generation; - } - wake_condition.notify_one(); - if (thread.joinable()) thread.join(); + poll_timer.cancel_and_wait(); } Gpu_Completion_Service::Gpu_Completion_Service() = default; @@ -145,11 +140,7 @@ Gpu_Completion_Service::Private::prepare( }); throw; } - { - std::lock_guard lock(service_mutex); - ++wake_generation; - } - wake_condition.notify_one(); + request_poll(std::chrono::nanoseconds{1}); return {Reservation(std::move(pending)), Admission_Result::none}; } @@ -166,7 +157,6 @@ void Gpu_Completion_Service::Private::watch( pending->fence = fence; pending->watched_at = std::chrono::steady_clock::now(); pending->status = Gpu_Completion_Pending_Fence::Status::watched; - ++wake_generation; } object->template update_state<&State::watched, &State::peak_watched>( [](State_Access states) { @@ -174,7 +164,7 @@ void Gpu_Completion_Service::Private::watch( ++state.watched; update_peak(state.peak_watched, state.watched); }); - wake_condition.notify_one(); + request_poll(std::chrono::nanoseconds{1}); } void Gpu_Completion_Service::Private::cancel( @@ -182,13 +172,10 @@ void Gpu_Completion_Service::Private::cancel( { std::lock_guard lock(service_mutex); if (!pending || pending->service != object) return; - if (pending->status == - Gpu_Completion_Pending_Fence::Status::reserved) - pending->status = - Gpu_Completion_Pending_Fence::Status::canceled; - ++wake_generation; + if (pending->status == Gpu_Completion_Pending_Fence::Status::reserved) + pending->status = Gpu_Completion_Pending_Fence::Status::canceled; } - wake_condition.notify_one(); + request_poll(std::chrono::nanoseconds{1}); } Gpu_Completion_State Gpu_Completion_Service::Private::state() const { @@ -220,17 +207,21 @@ Gpu_Completion_State Gpu_Completion_Service::Private::state() const { return result; } -void Gpu_Completion_Service::Private::run() noexcept { - struct Device_Fences { - VkDevice device{VK_NULL_HANDLE}; /* 同组 fence 的 Vulkan Device。 */ - std::vector fences{}; /* 本轮等待的同 Device fence 集合。 */ - }; - std::vector> active; +void Gpu_Completion_Service::Private::request_poll( + std::chrono::nanoseconds delay) noexcept { + if (stopping.load(std::memory_order_acquire)) return; + try { poll_timer.start_once(delay); } + catch (...) {} +} + +void Gpu_Completion_Service::Private::poll() noexcept { + if (stopping.load(std::memory_order_acquire)) return; try { - active.reserve(static_cast(default_capacity)); - std::size_t wait_group_index{}; - auto last_state_publication = - std::chrono::steady_clock::time_point{}; + object->template exchange_stream(); + object->template access_rendering_stream( + [&](std::span> pending) { + active.insert(active.end(), pending.begin(), pending.end()); + }); const auto finish = [this]( const std::shared_ptr& pending, @@ -240,22 +231,18 @@ void Gpu_Completion_Service::Private::run() noexcept { Result result{}; { std::lock_guard lock(service_mutex); - if (pending->status != - Gpu_Completion_Pending_Fence::Status::watched) + if (pending->status != Gpu_Completion_Pending_Fence::Status::watched) return; completion = std::move(pending->completion); on_exception = std::move(pending->on_exception); result.error = error; result.vulkan_result = vulkan_result; if (pending->observe) - result.wait_duration_ns = - elapsed_nanoseconds(pending->watched_at); - pending->status = - Gpu_Completion_Pending_Fence::Status::canceled; + result.wait_duration_ns = elapsed_nanoseconds(pending->watched_at); + pending->status = Gpu_Completion_Pending_Fence::Status::canceled; } object->template update_state< - &State::watched, &State::fault_count, - &State::abandoned_count>( + &State::watched, &State::fault_count, &State::abandoned_count>( [error](State_Access states) { auto& state = states.template get(); if (state.watched != 0) --state.watched; @@ -266,8 +253,7 @@ void Gpu_Completion_Service::Private::run() noexcept { } }); - const auto callback_started = - std::chrono::steady_clock::now(); + 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(); } @@ -287,183 +273,95 @@ void Gpu_Completion_Service::Private::run() noexcept { if (!callback_failure) return; try { on_exception(contextual_exception( - "delivering GPU completion", - std::move(callback_failure))); + "delivering GPU completion", std::move(callback_failure))); } catch (...) {} }; - for (;;) { - object->template exchange_stream(); - object->template access_rendering_stream( - [&](std::span> pending) { - active.insert(active.end(), pending.begin(), pending.end()); - }); - std::uint64_t generation{}; - std::size_t reserved_count{}; - std::vector groups; + std::size_t reserved_count{}; + bool needs_more{}; + 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{}; + Gpu_Completion_Pending_Fence::Status status{}; { std::lock_guard lock(service_mutex); - generation = wake_generation; - for (auto iterator = active.begin(); iterator != active.end();) { - const auto status = (*iterator)->status; - if (status == - Gpu_Completion_Pending_Fence::Status::canceled || - (status == - Gpu_Completion_Pending_Fence::Status::reserved && - stopping.load(std::memory_order_acquire))) { - iterator = active.erase(iterator); - object->template update_state< - &State::cancellation_count, &State::in_flight>( - [](State_Access states) { - auto& state = states.template get(); - ++state.cancellation_count; - if (state.in_flight != 0) --state.in_flight; - }); - continue; - } - if (status == - Gpu_Completion_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; - } + status = (*iterator)->status; + device = (*iterator)->device; + fence = (*iterator)->fence; + watched_at = (*iterator)->watched_at; } - - 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) { - object->template update_state< - &State::active_fences, &State::pending_fences>( - [active_count = active.size(), reserved_count]( - State_Access states) { - auto& state = states.template get(); - state.active_fences = active_count; - state.pending_fences = reserved_count; - }); - object->template publish_state(); - last_state_publication = publication_time; - } - if (stopping.load(std::memory_order_acquire) && active.empty()) { - object->template update_state< - &State::active_fences, &State::pending_fences>( - [](State_Access states) { - auto& state = states.template get(); - state.active_fences = 0; - state.pending_fences = 0; - }); - object->template publish_state(); - 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{}; - Gpu_Completion_Pending_Fence::Status status{}; - { - std::lock_guard lock(service_mutex); - status = (*iterator)->status; - device = (*iterator)->device; - fence = (*iterator)->fence; - watched_at = (*iterator)->watched_at; - } - if (status != - Gpu_Completion_Pending_Fence::Status::watched) { - ++iterator; - continue; - } - if (stopping.load(std::memory_order_acquire) || - 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; - } - object->template update_state<&State::fence_probe_count>( - [](State_Access states) { - ++states.template get().fence_probe_count; - }); - const VkResult result = vkGetFenceStatus(device, fence); - if (result == VK_NOT_READY) { - ++iterator; - continue; - } - auto pending = *iterator; + if (status == Gpu_Completion_Pending_Fence::Status::canceled) { 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, generation] { - return wake_generation != generation || - stopping.load(std::memory_order_acquire); - }); + object->template update_state< + &State::cancellation_count, &State::in_flight>( + [](State_Access states) { + auto& state = states.template get(); + ++state.cancellation_count; + if (state.in_flight != 0) --state.in_flight; + }); continue; } - - wait_group_index %= groups.size(); - const auto& group = groups[wait_group_index++]; - const auto wait_started = std::chrono::steady_clock::now(); - object->template update_state<&State::fence_wait_count>( + if (status == Gpu_Completion_Pending_Fence::Status::reserved) { + ++reserved_count; + ++iterator; + continue; + } + if (elapsed_nanoseconds(watched_at) >= maximum_fence_age_ns) { + auto pending = *iterator; + iterator = active.erase(iterator); + finish(pending, VK_TIMEOUT, Completion_Error::fence_abandoned); + continue; + } + object->template update_state<&State::fence_probe_count>( [](State_Access states) { - ++states.template get().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); - object->template update_state< - &State::fence_wait_total_ns, &State::fence_wait_max_ns, - &State::fence_wait_timeout_count>( - [wait_ns, wait_result](State_Access states) { - auto& state = states.template get(); - 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; + ++states.template get().fence_probe_count; }); + const VkResult result = vkGetFenceStatus(device, fence); + if (result == VK_NOT_READY) { + needs_more = true; + ++iterator; + continue; + } + auto pending = *iterator; + iterator = active.erase(iterator); + finish(pending, result, + result == VK_SUCCESS ? Completion_Error::none + : Completion_Error::vulkan_failure); } + + const auto now = std::chrono::steady_clock::now(); + if (last_state_publication == std::chrono::steady_clock::time_point{} || + now - last_state_publication >= state_publication_interval) { + object->template update_state< + &State::active_fences, &State::pending_fences>( + [active_count = active.size(), reserved_count]( + State_Access states) { + auto& state = states.template get(); + state.active_fences = active_count; + state.pending_fences = reserved_count; + }); + object->template publish_state(); + last_state_publication = now; + } + + if (needs_more) request_poll(probe_interval); } catch (...) { - const auto service_failure = std::current_exception(); - stopping.store(true, std::memory_order_release); - for (auto& pending : active) { - Exception_Handler handler; - { - std::lock_guard lock(service_mutex); - pending->status = - Gpu_Completion_Pending_Fence::Status::canceled; - handler = pending->on_exception; + const auto failure = std::current_exception(); + std::vector handlers; + { + std::lock_guard lock(service_mutex); + for (auto& pending : active) { + pending->status = Gpu_Completion_Pending_Fence::Status::canceled; + if (pending->on_exception) + handlers.push_back(std::move(pending->on_exception)); } - if (!handler) continue; - try { handler(service_failure); } + active.clear(); + } + for (auto& handler : handlers) { + try { handler(failure); } catch (...) {} } } diff --git a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp index 03ff3e3..ac22bce 100644 --- a/render_3D/render_3D/detail/Gpu_Completion_Service.hpp +++ b/render_3D/render_3D/detail/Gpu_Completion_Service.hpp @@ -19,7 +19,7 @@ struct Gpu_Completion_Service : Def #include #include -#include #include -#include +#include namespace aethera::render_3d::detail { struct Gpu_Completion_Pending_Fence { enum struct Status : std::uint8_t { @@ -23,11 +23,11 @@ struct Gpu_Completion_Pending_Fence { struct Gpu_Completion_Service::Private : Prev_Private { using Object = Impl; Object* object{}; /* Def 机制所属最终对象;生命周期与本 Private 相同。 */ - std::mutex service_mutex{}; /* 仅保护 reservation 状态与条件变量版本。 */ - std::condition_variable wake_condition{}; /* 专用完成线程的休眠唤醒源。 */ - std::uint64_t wake_generation{}; /* service_mutex 下的防丢唤醒版本。 */ + std::mutex service_mutex{}; /* 保护 reservation 状态;fence probe 只在全局 Scheduler 控制线程执行。 */ std::atomic_bool stopping{}; /* Private 析构开始后拒绝新准入。 */ - std::thread thread{}; /* 专用 fence 完成线程。 */ + std::vector> active{}; /* 仅 Frame_Scheduler 控制线程访问。 */ + Frame_Scheduler::Timer poll_timer{}; /* 与 2D/3D 帧时钟共用的唯一控制线程,不再单独创建 completion thread。 */ + std::chrono::steady_clock::time_point last_state_publication{}; Private(); ~Private(); template @@ -37,9 +37,9 @@ struct Gpu_Completion_Service::Private : Prev_Private { object->template update_state<&State::capacity>( static_cast(default_capacity)); object->template publish_state(); - thread = std::thread([this] { - run(); - }); + active.reserve(static_cast(default_capacity)); + poll_timer = Frame_Scheduler::instance().make_timer( + [this](Frame_Scheduler::Tick) { poll(); }); } [[nodiscard]] Prepare_Result prepare(Completion completion, Exception_Handler on_exception, @@ -48,11 +48,12 @@ struct Gpu_Completion_Service::Private : Prev_Private { VkDevice device, VkFence fence); void cancel(const std::shared_ptr& pending); [[nodiscard]] Gpu_Completion_State state() const; - void run() noexcept; + void request_poll(std::chrono::nanoseconds delay = probe_interval) noexcept; + void poll() noexcept; static constexpr std::ptrdiff_t default_capacity = 1024; - static constexpr std::uint64_t fence_wait_timeout_ns = 1'000'000; + static constexpr auto probe_interval = std::chrono::milliseconds(1); static constexpr std::uint64_t maximum_fence_age_ns = 30'000'000'000ULL; static constexpr auto state_publication_interval = - std::chrono::milliseconds(100); + std::chrono::milliseconds(100); }; } diff --git a/web_server/src/Gallery_Video_Stream.cpp b/web_server/src/Gallery_Video_Stream.cpp index d27e039..8d137bd 100644 --- a/web_server/src/Gallery_Video_Stream.cpp +++ b/web_server/src/Gallery_Video_Stream.cpp @@ -2,7 +2,6 @@ #include #include #include "detail/Gallery_Frame_Atlas.hpp" -#include "detail/Gallery_Frame_Clock.hpp" #include #include #include @@ -181,7 +180,6 @@ void Gallery_Video_Stream::Private::publish( void Gallery_Video_Stream::Private::fail( std::exception_ptr failure) noexcept { if (failed.exchange(true, std::memory_order_acq_rel)) return; - if (clock) clock->stop(); try { terminal_failure = exception_description(failure); publish(std::nullopt, nlohmann::json{ @@ -200,21 +198,13 @@ void Gallery_Video_Stream::Private::accept_frame( return; static_cast(atlas->accept_frame(slot, frame->pixels)); if (slot == encode_source_slot && frame->pixels) { + delivered_clock_ticks.fetch_add(1, std::memory_order_relaxed); active_encode_tick.sequence = frame->pixels->correlation_id; active_encode_tick.time_milliseconds = static_cast(frame->pixels->presentation_time.count()) / 1'000.0; } } -void Gallery_Video_Stream::Private::dispatch_tick(Plot_Render_Tick tick) { - if (stopping.load(std::memory_order_acquire) || - failed.load(std::memory_order_acquire)) - return; - tick.width = tile_width; - tick.height = tile_height; - delivered_clock_ticks.fetch_add(1, std::memory_order_relaxed); - for (const auto& source : sources) source.entry.plot->schedule_render(tick); -} void Gallery_Video_Stream::Private::update_metrics( const Plot_Render_Tick& tick, const detail::Gallery_Atlas_Composition& composition, @@ -395,22 +385,6 @@ Gallery_Video_Stream::~Gallery_Video_Stream() { void Gallery_Video_Stream::bind_plots() { auto& data = static_cast(*d); const auto weak = weak_from_this(); - data.clock = std::make_unique( - gallery_frame_rate, - [weak](Plot_Render_Tick tick) { - if (const auto owner = weak.lock()) static_cast(*owner->d).dispatch_tick(tick); - }, - [weak](std::exception_ptr failure) { - if (const auto owner = weak.lock()) { - aethera::schedule_task( - "web.gallery.failure", - [weak, failure = std::move(failure)]() mutable { - if (const auto stream = weak.lock()) - static_cast(*stream->d).fail( - std::move(failure)); - }); - } - }); for (std::size_t slot = 0; slot < data.sources.size(); ++slot) { auto& source = data.sources[slot]; source.stream = source.entry.plot->subscribe( @@ -515,13 +489,6 @@ Gallery_Video_Stream::Stream_Id Gallery_Video_Stream::subscribe( std::memory_order_acquire)) break; } - try { - data.clock->start(); - } - catch (...) { - static_cast(data.remove_consumer(id)); - throw; - } return id; } void Gallery_Video_Stream::unsubscribe(Stream_Id stream) { @@ -537,7 +504,6 @@ void Gallery_Video_Stream::request_video_key_frame() { void Gallery_Video_Stream::shutdown() noexcept { auto& data = static_cast(*d); if (data.stopping.exchange(true, std::memory_order_acq_rel)) return; - if (data.clock) data.clock->stop(); for (auto& source : data.sources) { if (source.stream == 0) continue; try { diff --git a/web_server/src/Gallery_Video_Stream.ipp b/web_server/src/Gallery_Video_Stream.ipp index cb99b14..099ed04 100644 --- a/web_server/src/Gallery_Video_Stream.ipp +++ b/web_server/src/Gallery_Video_Stream.ipp @@ -1,6 +1,5 @@ #pragma once #include "detail/Gallery_Frame_Atlas.hpp" -#include "detail/Gallery_Frame_Clock.hpp" #include #include #include @@ -28,7 +27,6 @@ struct Gallery_Video_Stream::Private : Prev_Private { Sliding_Statistics publish_ms{600}; /* WebRTC 发布统计计算器。 */ std::vector sources{}; /* 已按业务标识排序的稳定图集来源。 */ std::unique_ptr atlas{}; /* 最近完成帧与 RGBA 图集的唯一状态源。 */ - std::unique_ptr clock{}; /* 页面全部 Plot 共用的绝对帧时钟。 */ H264_Encoder encoder; /* 页面唯一 H.264 编码器。 */ std::unique_ptr encode_domain{}; /* FFmpeg 驱动调用域。 */ Plot_Render_Tick active_encode_tick{}; /* 当前媒体 DAG 的输入时钟。 */ @@ -65,7 +63,6 @@ struct Gallery_Video_Stream::Private : Prev_Private { void fail(std::exception_ptr failure) noexcept; void accept_frame(std::size_t slot, std::shared_ptr frame); - void dispatch_tick(Plot_Render_Tick tick); void update_metrics(const Plot_Render_Tick& tick, const detail::Gallery_Atlas_Composition& composition, std::size_t encoded_bytes, diff --git a/web_server/src/Plot.cpp b/web_server/src/Plot.cpp index 1607377..ed8dc7e 100644 --- a/web_server/src/Plot.cpp +++ b/web_server/src/Plot.cpp @@ -1,6 +1,8 @@ #include "Plot.hpp" #include "Renderable_Adapter.hpp" #include "Taskflow_Trace_Json.hpp" +#include +#include #include #include #include @@ -47,30 +49,22 @@ std::string exception_description(const std::exception_ptr& failure) { return "empty Plot failure"; } -enum struct Frame_Pacing_Mode : std::uint8_t { - manual, - fixed_rate, - maximum_rate -}; - -struct Frame_Pacing_Properties { - bool render_enabled{true}; - bool video_enabled{true}; - Frame_Pacing_Mode mode{Frame_Pacing_Mode::fixed_rate}; /* 只决定共享时钟到达时是否调用 Scene::render。 */ - double fixed_rate_fps{100.0}; /* 页面级时钟当前上限为每秒一百次。 */ -}; - struct Frame_Policy final { public: - [[nodiscard]] Frame_Pacing_Properties read() const { - return *pacing.load(std::memory_order_acquire); - } + [[nodiscard]] Frame_Pacing_Properties read() const { return pacing.read(); } [[nodiscard]] nlohmann::json schema() const; [[nodiscard]] nlohmann::json write_prop(std::string_view key, const nlohmann::json& value); + [[nodiscard]] bool accept_periodic_tick(double time_milliseconds) { + return pacing.accept_periodic_tick(time_milliseconds); + } + [[nodiscard]] bool request_immediate() { return pacing.request_immediate(); } + [[nodiscard]] double scheduled_rate_fps() const { return pacing.scheduled_rate_fps(); } + void frame_submitted() { pacing.frame_submitted(); } + [[nodiscard]] bool frame_completed() { return pacing.frame_completed(); } + void frame_rejected() { pacing.frame_rejected(); } private: - std::atomic> pacing{ - std::make_shared()}; + Frame_Pacing_Policy pacing{}; }; std::string_view pacing_mode_name(Frame_Pacing_Mode mode) { @@ -103,26 +97,26 @@ nlohmann::json Frame_Policy::schema() const { {"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."}, + {"editable", true}, {"description", "控制当前 Scene 的周期刷新;画面隐藏不会修改此项。"}, + {"technical_description", "Authoritative per-scene periodic render switch."}, {"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", 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."}, + {"editable", true}, {"description", "只控制 Scene::render(Frame*) 的调用节奏;Scene 的 Frame 所有权与接口保持不变。"}, + {"technical_description", "Per-scene frame pacing policy backed by the Kernel scheduler."}, {"value", pacing_mode_name(current.mode)}, {"options", nlohmann::json::array({ {{"value", "manual"}, {"label", "手动渲染"}}, {{"value", "fixed_rate"}, {"label", "固定频率"}}, - {{"value", "maximum_rate"}, {"label", "共享时钟满速"}} + {{"value", "maximum_rate"}, {"label", "最大频率"}} })}}, {{"key", "fixed_rate_fps"}, {"label", "目标帧率"}, {"editor", "number"}, {"editable", true}, {"minimum", 0.1}, {"maximum", 100.0}, {"step", 0.1}, - {"description", "固定频率模式下从页面级 100 Hz 时钟采样当前 Scene 的次数。"}, - {"technical_description", "Per-plot rate selected from the common gallery render timeline."}, + {"description", "当前 Scene 独立目标帧率;周期策略保存在 Kernel Frame_Pacing_Policy。"}, + {"technical_description", "Independent per-scene target frame rate."}, {"value", current.fixed_rate_fps}} })} }; @@ -130,26 +124,12 @@ nlohmann::json Frame_Policy::schema() const { nlohmann::json Frame_Policy::write_prop(std::string_view key, const nlohmann::json& value) { - 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"}}; const bool target = value.get(); - update([&](Frame_Pacing_Properties& next) { - (key == "render_enabled" ? next.render_enabled : next.video_enabled) = - target; - }); + if (key == "render_enabled") pacing.set_render_enabled(target); + else pacing.set_video_enabled(target); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, {"value", target}}; } @@ -159,7 +139,7 @@ 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"}}; - update([&](Frame_Pacing_Properties& next) { next.mode = *parsed; }); + pacing.set_mode(*parsed); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, {"value", pacing_mode_name(*parsed)}}; } @@ -169,9 +149,7 @@ 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"}}; - update([&](Frame_Pacing_Properties& properties) { - properties.fixed_rate_fps = next; - }); + pacing.set_fixed_rate(next); return {{"success", true}, {"component", "frame-analysis"}, {"key", key}, {"value", next}}; } @@ -363,19 +341,20 @@ struct Plot::Private { std::unique_ptr view; std::once_flag start_once; + std::weak_ptr lifetime{}; /* 仅用于 completion 后重新投递 Taskflow,避免在 Scene callback 内重入 render。 */ 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{}; + Frame_Scheduler::Timer frame_timer{}; /* 每 Plot/Scene 只有轻量时间轮节点,不持有线程。 */ static constexpr std::size_t scene_frame_capacity{3}; std::array frame_slots{}; /* Scene 借用的稳定三缓冲物理帧。 */ std::atomic_size_t callback_retired_slot{scene_frame_capacity}; Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */ 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 总数。 */ std::atomic_uint64_t coalesced_tick_count{}; /* 尚未消费时被更新 tick 替换的旧 tick 总数。 */ @@ -419,6 +398,7 @@ struct Plot::Private { [[nodiscard]] Stream_Snapshot stream_snapshot() const; void publish(std::shared_ptr frame) noexcept; void consume_tick(std::weak_ptr lifetime); + void refresh_schedule(); void clock_tick(const Plot_Render_Tick& tick); void render_frame(Plot_Render_Tick tick); void publish_completed_frame(); @@ -501,6 +481,14 @@ void Plot::Private::publish( catch (...) {} } +void Plot::Private::refresh_schedule() { + if (!frame_timer.valid()) return; + const auto current_consumers = consumers.load(std::memory_order_acquire); + const double fps = frame_policy.scheduled_rate_fps(); + if (fps <= 0.0 || current_consumers->empty()) frame_timer.cancel(); + else frame_timer.start_periodic(fps); +} + void Plot::Private::consume_tick(std::weak_ptr lifetime) { const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel); if (tick) clock_tick(*tick); @@ -519,21 +507,10 @@ 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.read(); - if (!pacing.render_enabled || pacing.mode == Frame_Pacing_Mode::manual) { + if (!frame_policy.accept_periodic_tick(tick.time_milliseconds)) { policy_skip_count.fetch_add(1, std::memory_order_relaxed); return; } - if (pacing.mode == Frame_Pacing_Mode::fixed_rate) { - if (tick.time_milliseconds < last_clock_render_time_ms) - last_clock_render_time_ms = -std::numeric_limits::infinity(); - const double interval = 1'000.0 / pacing.fixed_rate_fps; - if (tick.time_milliseconds - last_clock_render_time_ms + 0.01 < interval) { - policy_skip_count.fetch_add(1, std::memory_order_relaxed); - return; - } - } - last_clock_render_time_ms = tick.time_milliseconds; render_frame(tick); } @@ -638,8 +615,10 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { restore_taskflow_trace_claim(); } else { taskflow_trace_claimed = false; + frame_policy.frame_submitted(); submitted_frame_count.fetch_add(1, std::memory_order_relaxed); } + if (!result) frame_policy.frame_rejected(); return; } auto& output = *std::get>(managed->frame); @@ -653,9 +632,11 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) { const auto result = scene_3d->render(&output); if (result == Render_Scene_3D::Render_Result::submitted) { taskflow_trace_claimed = false; + frame_policy.frame_submitted(); submitted_frame_count.fetch_add(1, std::memory_order_relaxed); return; } + frame_policy.frame_rejected(); scene_rejection_count.fetch_add(1, std::memory_order_relaxed); rollback_unsubmitted(); restore_taskflow_trace_claim(); @@ -794,6 +775,23 @@ void Plot::Private::retire_completed_frame(Render_Frame* frame) { break; } } + if (frame_policy.frame_completed()) { + const auto weak = lifetime; + aethera::schedule_task("web.plot.render.deferred", [weak] { + const auto owner = weak.lock(); + if (!owner) return; + try { + const auto now = std::chrono::steady_clock::now(); + const auto elapsed = now - owner->d->clock_origin; + owner->d->render_frame(Plot_Render_Tick{ + now, 0, + std::chrono::duration(elapsed).count()}); + } + catch (...) { + owner->d->fail(std::current_exception()); + } + }); + } } void Plot::Private::attach_completion( @@ -828,14 +826,24 @@ void Plot::attach_scene_completion( void Plot::ensure_started() { std::call_once(d->start_once, [this] { + const auto weak = weak_from_this(); + d->lifetime = weak; + d->frame_timer = Frame_Scheduler::instance().make_timer( + [weak](Frame_Scheduler::Tick tick) { + if (const auto owner = weak.lock()) { + owner->schedule_render(Plot_Render_Tick{ + tick.issued_at, tick.sequence, tick.time_milliseconds}); + } + }); + d->refresh_schedule(); /* - * 页面级帧时钟只触发 render;Scene 完成回调直接归还帧并发布结果。 - * 回调不请求下一帧,也不重新投递任务;Scene 在回调返回时释放准入, - * 提交和页面级媒体图集能够按各自的资源域自然流水。 + * Kernel Frame_Scheduler 只产生每个 Scene 独立的周期 tick; + * Scene::render(Frame*) 与外部 Frame 所有权保持原样。完成回调先归还 + * 当前帧;若期间出现 immediate 请求,则回调返回后重新投递 Taskflow。 */ if (auto* scene = std::get_if>(&d->scene)) { (*scene)->set_frame_callback( - [weak = weak_from_this()](Frame_2D* frame) { + [weak](Frame_2D* frame) { if (auto owner = weak.lock()) { try { owner->d->retire_completed_frame(frame); } catch (...) { owner->d->fail(std::current_exception()); } @@ -843,7 +851,7 @@ void Plot::ensure_started() { }); } else { std::get>(d->scene)->set_frame_callback( - [weak = weak_from_this()](Frame_3D* frame) { + [weak](Frame_3D* frame) { if (auto owner = weak.lock()) { try { owner->d->retire_completed_frame(frame); } catch (...) { owner->d->fail(std::current_exception()); } @@ -869,6 +877,7 @@ Plot::Stream_Id Plot::subscribe(Stream_Handler handler) { std::memory_order_acquire)) break; } + d->refresh_schedule(); if (d->terminal_failure.load(std::memory_order_acquire)) { const auto failure = d->terminal_failure.load(std::memory_order_acquire); try { @@ -896,8 +905,9 @@ void Plot::unsubscribe(Stream_Id stream) { if (d->consumers.compare_exchange_weak( current, desired, std::memory_order_release, std::memory_order_acquire)) - return; + break; } + d->refresh_schedule(); } void Plot::configure_stream(Stream_Id stream, std::uint32_t width, @@ -941,6 +951,7 @@ void Plot::schedule_render(Plot_Render_Tick tick) { void Plot::render_once() { ensure_started(); + if (!d->frame_policy.request_immediate()) return; const auto weak = weak_from_this(); aethera::schedule_task("web.plot.render.once", [weak] { const auto owner = weak.lock(); @@ -987,9 +998,11 @@ nlohmann::json Plot::write_prop(std::string_view component, std::string_view key, const nlohmann::json& value) { ensure_started(); - return component == "frame-analysis" - ? d->frame_policy.write_prop(key, value) - : d->view->write_prop(component, key, value); + if (component != "frame-analysis") + return d->view->write_prop(component, key, value); + auto result = d->frame_policy.write_prop(key, value); + if (result.value("success", false)) d->refresh_schedule(); + return result; } nlohmann::json Plot::generate_data(const nlohmann::json& input) { diff --git a/web_server/src/Plot.hpp b/web_server/src/Plot.hpp index 16d76e6..8454123 100644 --- a/web_server/src/Plot.hpp +++ b/web_server/src/Plot.hpp @@ -29,7 +29,7 @@ struct Plot_Input_Event { }; struct Plot_Render_Tick { std::chrono::steady_clock::time_point issued_at{}; /* 页面帧时钟发布本 tick 的单调时刻;手动帧在提交时填写。 */ - std::uint64_t sequence{}; /* 页面级帧时钟分配的关联序号。 */ + std::uint64_t sequence{}; /* Kernel 全局 Frame_Scheduler 时间轴上的关联序号。 */ double time_milliseconds{}; /* 页面级单调时间线,所有图共享同一个动画时刻。 */ std::uint32_t width{320}; /* 当前图在媒体图集中的固定像素宽度。 */ std::uint32_t height{192}; /* 当前图在媒体图集中的固定像素高度。 */ diff --git a/web_server/src/detail/Gallery_Frame_Clock.cpp b/web_server/src/detail/Gallery_Frame_Clock.cpp deleted file mode 100644 index 5a52d0f..0000000 --- a/web_server/src/detail/Gallery_Frame_Clock.cpp +++ /dev/null @@ -1,109 +0,0 @@ -#include "Gallery_Frame_Clock.hpp" -#include -#include -#include -#include -#include -#include -#include -#include - -namespace aethera::web::detail { -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) - : tick_handler(std::move(value_tick_handler)), - failure_handler(std::move(value_failure_handler)) { - if (!std::isfinite(frame_rate) || frame_rate <= 0.0) - throw std::invalid_argument("gallery frame rate must be positive"); - interval = std::chrono::nanoseconds{ - static_cast(std::llround(1'000'000'000.0 / frame_rate))}; - if (interval.count() <= 0) - throw std::invalid_argument("gallery frame rate exceeds clock resolution"); - if (!tick_handler) - throw std::invalid_argument("gallery frame clock requires a tick handler"); - } - - 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, std::uint64_t active_generation) noexcept { - const auto origin = std::chrono::steady_clock::now(); - std::uint64_t sequence{}; - auto deadline = origin + interval; - while (!stop.stop_requested()) { - std::unique_lock lock(lifecycle_mutex); - wake.wait_until(lock, stop, deadline, [] { return false; }); - lock.unlock(); - if (stop.stop_requested()) break; - try { - ++sequence; - tick_handler(Plot_Render_Tick{ - std::chrono::steady_clock::now(), - sequence, - std::chrono::duration(deadline - origin).count()}); - deadline += interval; - const auto now = std::chrono::steady_clock::now(); - if (deadline <= now) - deadline += interval * ((now - deadline) / interval + 1); - } - catch (...) { - report(active_generation, std::current_exception()); - return; - } - } - if (generation.load(std::memory_order_acquire) == active_generation) - running.store(false, std::memory_order_release); - } -}; - -Gallery_Frame_Clock::Gallery_Frame_Clock( - double frame_rate, Tick_Handler tick_handler, - Failure_Handler failure_handler) - : d(std::make_shared(frame_rate, std::move(tick_handler), - std::move(failure_handler))) {} - -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; - 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) 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/src/detail/Gallery_Frame_Clock.hpp b/web_server/src/detail/Gallery_Frame_Clock.hpp deleted file mode 100644 index d1c1ed1..0000000 --- a/web_server/src/detail/Gallery_Frame_Clock.hpp +++ /dev/null @@ -1,22 +0,0 @@ -#pragma once -#include "../Plot.hpp" -#include -#include -#include -namespace aethera::web::detail { -struct Gallery_Frame_Clock final { -public: - using Tick_Handler = std::function; - using Failure_Handler = std::function; - Gallery_Frame_Clock(double frame_rate, Tick_Handler tick_handler, - Failure_Handler failure_handler); - ~Gallery_Frame_Clock(); - Gallery_Frame_Clock(const Gallery_Frame_Clock&) = delete; - Gallery_Frame_Clock& operator=(const Gallery_Frame_Clock&) = delete; - void start(); - void stop() noexcept; -private: - struct Private; - std::shared_ptr d; /* 异步计时操作完成前保持时钟状态存活。 */ -}; -}