加上纯色websocket

This commit is contained in:
2026-08-26 18:06:48 +08:00
parent 272d85caac
commit 2de8a35fff
28 changed files with 2061 additions and 549 deletions
+370 -180
View File
@@ -8,6 +8,7 @@
#include <render_2D/plottable/Plottables.hpp>
#include <render_3D/Render_3D.hpp>
#include <render_3D/Gpu_Completion_State.hpp>
#include <render_common.hpp>
#include <algorithm>
#include <array>
#include <atomic>
@@ -18,6 +19,7 @@
#include <initializer_list>
#include <limits>
#include <memory>
#include <mutex>
#include <optional>
#include <span>
#include <stdexcept>
@@ -101,7 +103,7 @@ nlohmann::json Frame_Policy::schema() const {
{"technical_description", "Authoritative per-scene periodic render switch."},
{"value", current.render_enabled}},
{{"key", "video_enabled"}, {"label", "图集视频传输"}, {"editor", "boolean"},
{"editable", true}, {"description", "控制完成帧是否进入页面级 RGBA 图集;默认开启,用于完整链路压测"},
{"editable", true}, {"description", "控制完成帧是否进入页面级采样器;2D BGRA 与 3D RGBA 均保持原生格式"},
{"technical_description", "Authoritative tile publication switch for the shared gallery video."},
{"value", current.video_enabled}},
{{"key", "pacing_mode"}, {"label", "服务端帧策略"}, {"editor", "select"},
@@ -319,8 +321,8 @@ struct Plot::Private {
enum struct Frame_State : std::uint8_t {
available,
in_flight,
callback_retired
rendering, /* Scene::advance -> Plot pixel publish,不可重入。 */
consuming /* 外接 Taskflow 正在消费已发布帧;允许下一帧渲染。 */
};
struct Managed_Frame {
@@ -353,16 +355,17 @@ struct Plot::Private {
Frame_Policy frame_policy{};
Frame_Scheduler::Timer frame_timer{}; /* 每 Plot/Scene 只有轻量时间轮节点,不持有线程。 */
static constexpr std::size_t scene_frame_capacity{3};
std::array<Managed_Frame, scene_frame_capacity> frame_slots{}; /* Scene 借用的稳定三缓冲物理帧。 */
std::atomic_size_t callback_retired_slot{scene_frame_capacity};
std::array<Managed_Frame, scene_frame_capacity> frame_slots{}; /* Scene 与外接消费者共享生命周期的稳定三缓冲。 */
Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */
std::atomic<std::shared_ptr<const Plot_Render_Tick>> pending_tick{};
std::atomic_bool tick_task_scheduled{}; /* 唯一短任务准入;不占用 Worker 等待。 */
std::atomic_bool render_admission_busy{}; /* view->update/Scene::advance 到 pixel publish 的唯一准入门。 */
std::chrono::steady_clock::time_point clock_origin{std::chrono::steady_clock::now()};
std::atomic_uint64_t received_tick_count{}; /* 页面时钟交付给本 Plot 的 tick 总数。 */
std::atomic_uint64_t coalesced_tick_count{}; /* 尚未消费时被更新 tick 替换的旧 tick 总数。 */
std::atomic_uint64_t policy_skip_count{}; /* 帧策略拒绝的 tick 总数。 */
std::atomic_uint64_t preparation_busy_count{}; /* 本 Plot Prepare 已在执行而跳过的提交次数。 */
std::atomic_uint64_t preparation_busy_count{}; /* Render admission 忙时被合并为 latest pending 的 tick。 */
std::atomic_uint64_t deferred_resume_count{}; /* pixel publish 后立即唤醒 latest pending 的次数。 */
std::atomic_uint64_t frame_slot_busy_count{}; /* 三个物理帧槽均被占用的提交次数。 */
std::atomic_uint64_t scene_rejection_count{}; /* Scene 单帧准入拒绝的提交次数。 */
std::atomic_uint64_t submitted_frame_count{}; /* 成功提交给 Scene 的帧总数。 */
@@ -373,8 +376,19 @@ struct Plot::Private {
std::atomic_uint64_t taskflow_trace_control{};
std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames> taskflow_trace_slots{};
Task_Node completion_tail{}; /* Scene 完成图中当前最后一个业务阶段。 */
std::vector<std::unique_ptr<Task_Graph>> completion_extensions{}; /* 生命周期覆盖 Scene 对子图的借用。 */
std::atomic_size_t post_publish_trace_remaining{}; /* 仅捕获 publish 后外接 DAG 的剩余样本。 */
std::atomic_uint64_t post_publish_trace_control{}; /* 高 32 位 requested,低 32 位 captured。 */
std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames> post_publish_trace_slots{};
Task_Node completion_tail{}; /* Scene 图内固定停在 plot.frame.publish。 */
Task_Graph post_publish_graph{"plot.post_publish"}; /* publish 后外接 DAG;不再占用 Scene render admission。 */
Task_Node post_publish_tail{};
bool has_post_publish_tail{};
std::atomic_bool post_publish_busy{};
std::vector<std::unique_ptr<Task_Graph>> completion_extensions{}; /* 生命周期覆盖 post-publish module 借用。 */
Frame_Statistics_Accumulator completed_frame_statistics{diagnostic_window_capacity};
Frame_Statistics_State completed_frame_statistics_state{};
mutable std::mutex completed_frame_statistics_mutex{};
template <typename Scene_Object>
Private(std::unique_ptr<Scene_Object> value_scene,
@@ -400,14 +414,28 @@ struct Plot::Private {
[[nodiscard]] nlohmann::json schema() const;
[[nodiscard]] Stream_Snapshot stream_snapshot() const;
void publish(std::shared_ptr<const Plot_Stream_Frame> frame) noexcept;
void defer_tick(const Plot_Render_Tick& tick);
void arm_tick_consumer(std::weak_ptr<Plot> lifetime);
void release_render_admission(std::weak_ptr<Plot> lifetime);
void consume_tick(std::weak_ptr<Plot> lifetime);
void refresh_schedule();
void clock_tick(const Plot_Render_Tick& tick);
void render_frame(Plot_Render_Tick tick);
void publish_completed_frame();
void consume_completed_frame(Render_Frame* frame);
void retire_completed_frame(Render_Frame* frame);
void attach_completion(std::unique_ptr<Task_Graph> completion);
[[nodiscard]] bool mark_taskflow_trace(Render_Frame& frame);
[[nodiscard]] bool mark_post_publish_taskflow_trace();
void store_trace(std::atomic_uint64_t& control,
std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames>& slots,
const Taskflow_Frame_Trace& trace);
[[nodiscard]] nlohmann::json trace_response(
const std::atomic_uint64_t& control,
const std::atomic_size_t& remaining,
const std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames>& slots) const;
void fail(std::exception_ptr failure) noexcept;
};
@@ -492,25 +520,64 @@ void Plot::Private::refresh_schedule() {
else frame_timer.start_periodic(fps);
}
void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel);
if (tick) clock_tick(*tick);
tick_task_scheduled.store(false, std::memory_order_release);
if (pending_tick.load(std::memory_order_acquire) &&
!tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) {
aethera::schedule_task("web.plot.tick.consume", [lifetime] {
const auto plot = lifetime.lock();
if (!plot) return;
try { plot->d->consume_tick(lifetime); }
catch (...) { plot->d->fail(std::current_exception()); }
});
void Plot::Private::defer_tick(const Plot_Render_Tick& tick) {
const auto next = std::make_shared<const Plot_Render_Tick>(tick);
auto current = pending_tick.load(std::memory_order_acquire);
for (;;) {
if (current) {
const bool current_immediate = current->sequence == 0;
const bool next_immediate = tick.sequence == 0;
if ((current_immediate && !next_immediate) ||
(current_immediate == next_immediate &&
current->issued_at >= tick.issued_at))
return;
}
if (pending_tick.compare_exchange_weak(
current, next, std::memory_order_acq_rel,
std::memory_order_acquire)) {
if (current)
coalesced_tick_count.fetch_add(1, std::memory_order_relaxed);
return;
}
}
}
void Plot::Private::arm_tick_consumer(std::weak_ptr<Plot> lifetime) {
if (terminal_failure.load(std::memory_order_acquire) ||
render_admission_busy.load(std::memory_order_acquire) ||
!pending_tick.load(std::memory_order_acquire))
return;
if (tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return;
aethera::schedule_task("web.plot.tick.consume", [lifetime] {
const auto plot = lifetime.lock();
if (!plot) return;
try { plot->d->consume_tick(lifetime); }
catch (...) { plot->d->fail(std::current_exception()); }
});
}
void Plot::Private::release_render_admission(std::weak_ptr<Plot> lifetime) {
if (!render_admission_busy.exchange(false, std::memory_order_acq_rel))
return;
if (pending_tick.load(std::memory_order_acquire))
deferred_resume_count.fetch_add(1, std::memory_order_relaxed);
arm_tick_consumer(std::move(lifetime));
}
void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
if (!render_admission_busy.load(std::memory_order_acquire)) {
const auto tick = pending_tick.exchange({}, std::memory_order_acq_rel);
if (tick) clock_tick(*tick);
}
tick_task_scheduled.store(false, std::memory_order_release);
arm_tick_consumer(std::move(lifetime));
}
void Plot::Private::clock_tick(const Plot_Render_Tick& tick) {
if (terminal_failure.load(std::memory_order_acquire)) return;
if (!frame_policy.accept_periodic_tick(tick.time_milliseconds)) {
if (tick.sequence != 0 &&
!frame_policy.accept_periodic_tick(tick.time_milliseconds)) {
policy_skip_count.fetch_add(1, std::memory_order_relaxed);
return;
}
@@ -530,27 +597,85 @@ bool Plot::Private::mark_taskflow_trace(Render_Frame& frame) {
return false;
}
bool Plot::Private::mark_post_publish_taskflow_trace() {
auto remaining = post_publish_trace_remaining.load(std::memory_order_acquire);
while (remaining != 0) {
if (post_publish_trace_remaining.compare_exchange_weak(
remaining, remaining - 1, std::memory_order_acq_rel,
std::memory_order_acquire))
return true;
}
return false;
}
void Plot::Private::store_trace(
std::atomic_uint64_t& control,
std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames>& slots,
const Taskflow_Frame_Trace& value) {
auto trace = std::make_shared<const Taskflow_Frame_Trace>(value);
auto state = control.load(std::memory_order_acquire);
for (;;) {
const auto requested = static_cast<std::uint32_t>(state >> 32U);
const auto captured = static_cast<std::uint32_t>(state);
if (captured >= requested) return;
slots[captured].store(trace, std::memory_order_release);
const auto next = (static_cast<std::uint64_t>(requested) << 32U) |
static_cast<std::uint64_t>(captured + 1U);
if (control.compare_exchange_weak(
state, next, std::memory_order_release,
std::memory_order_acquire))
return;
}
}
nlohmann::json Plot::Private::trace_response(
const std::atomic_uint64_t& control,
const std::atomic_size_t& remaining,
const std::array<std::atomic<std::shared_ptr<const Taskflow_Frame_Trace>>,
maximum_taskflow_trace_frames>& slots) const {
nlohmann::json frames = nlohmann::json::array();
const auto state = control.load(std::memory_order_acquire);
const auto requested = static_cast<std::uint32_t>(state >> 32U);
const auto captured = static_cast<std::uint32_t>(state);
for (std::uint32_t index = 0; index < captured; ++index)
if (const auto trace = slots[index].load(std::memory_order_acquire))
frames.push_back(taskflow_trace_json(*trace));
const auto left = remaining.load(std::memory_order_acquire);
return {
{"protocol", "aethera.taskflow.frames"}, {"version", 1},
{"requested", requested}, {"remaining", left},
{"captured", frames.size()},
{"complete", requested != 0 && frames.size() == requested},
{"frames", std::move(frames)}};
}
void Plot::Private::render_frame(Plot_Render_Tick tick) {
if (terminal_failure.load(std::memory_order_acquire)) return;
const auto streams = stream_snapshot();
const auto pacing = frame_policy.read();
if (!pacing.render_enabled || streams.consumers->empty()) return;
std::size_t slot_index{};
Managed_Frame* managed{};
/* Frame_State 是 Plot 数据准备生命周期的唯一权威来源。槽位通过 CAS
* 准入,完成图和时钟任务不再共享一把外围锁。 */
if (std::ranges::any_of(frame_slots, [](const Managed_Frame& slot) {
return slot.state.load(std::memory_order_acquire) ==
Frame_State::in_flight;
})) {
bool admission_expected = false;
if (!render_admission_busy.compare_exchange_strong(
admission_expected, true, std::memory_order_acq_rel,
std::memory_order_acquire)) {
preparation_busy_count.fetch_add(1, std::memory_order_relaxed);
defer_tick(tick);
return;
}
std::size_t slot_index{};
Managed_Frame* managed{};
/*
* 只有 rendering 槽受 Scene 不可重入门约束;consuming 槽表示上一帧
* 已经完成 Plot 像素发布,外接 H264/WebRTC 仍可继续持有该物理帧的
* 诊断生命周期。只要还有 available 槽,下一帧即可进入。
*/
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
auto expected = Frame_State::available;
if (!frame_slots[index].state.compare_exchange_strong(
expected, Frame_State::in_flight,
expected, Frame_State::rendering,
std::memory_order_acq_rel, std::memory_order_acquire))
continue;
slot_index = index;
@@ -559,6 +684,14 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
}
if (!managed) {
frame_slot_busy_count.fetch_add(1, std::memory_order_relaxed);
defer_tick(tick);
/*
* 三个槽都仍被外接消费者持有时,只保留 latest pending。这里绝不能
* 立即 arm tick consumer,否则会在没有任何槽可用期间形成
* consume -> no slot -> consume 的 Taskflow 任务风暴。真正的唤醒点
* 是 retire_completed_frame:某个 consuming 槽变回 available 后只唤醒一次。
*/
render_admission_busy.store(false, std::memory_order_release);
return;
}
managed->presentation_time =
@@ -566,7 +699,7 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
std::chrono::duration<double, std::milli>(tick.time_milliseconds));
const auto rollback_unsubmitted = [this, slot_index] {
auto& slot = frame_slots[slot_index];
auto expected = Frame_State::in_flight;
auto expected = Frame_State::rendering;
static_cast<void>(slot.state.compare_exchange_strong(
expected, Frame_State::available, std::memory_order_acq_rel,
std::memory_order_acquire));
@@ -616,6 +749,7 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
scene_rejection_count.fetch_add(1, std::memory_order_relaxed);
rollback_unsubmitted();
restore_taskflow_trace_claim();
release_render_admission(lifetime);
} else {
taskflow_trace_claimed = false;
frame_policy.frame_submitted();
@@ -643,12 +777,14 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
scene_rejection_count.fetch_add(1, std::memory_order_relaxed);
rollback_unsubmitted();
restore_taskflow_trace_claim();
release_render_admission(lifetime);
if (result == Render_Scene_3D::Render_Result::backend_unavailable)
throw std::runtime_error("3D render backend became unavailable before submission");
}
catch (...) {
rollback_unsubmitted();
restore_taskflow_trace_claim();
release_render_admission(lifetime);
throw;
}
}
@@ -657,29 +793,19 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
void Plot::Private::publish_completed_frame() {
Render_Frame* frame{};
Managed_Frame* managed{};
/* Scene 在用户回调返回之后才标记 callback_finished。下一次串行完成图
* 开始时,上一 callback_retired 槽才可归还。 */
const auto retired_index = callback_retired_slot.exchange(
scene_frame_capacity, std::memory_order_acq_rel);
if (retired_index != scene_frame_capacity) {
auto expected = Frame_State::callback_retired;
if (!frame_slots[retired_index].state.compare_exchange_strong(
expected, Frame_State::available, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("Plot retired frame state is inconsistent");
}
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
if (frame_slots[index].state.load(std::memory_order_acquire) !=
Frame_State::in_flight) continue;
Frame_State::rendering)
continue;
if (managed)
throw std::logic_error("Plot has multiple active Scene frames");
throw std::logic_error("Plot has multiple frames in Scene rendering");
frame = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
managed = &frame_slots[index];
}
if (!managed)
throw std::logic_error("Scene completion graph has no active Plot frame");
throw std::logic_error("Scene completion graph has no rendering Plot frame");
try {
const auto pacing = frame_policy.read();
@@ -729,89 +855,144 @@ void Plot::Private::publish_completed_frame() {
static_cast<std::uint64_t>(std::max<std::int64_t>(0,
std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::steady_clock::now() - publish_started).count())));
auto expected = Frame_State::rendering;
if (!managed->state.compare_exchange_strong(
expected, Frame_State::consuming, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("Plot frame left rendering before pixel publish");
}
catch (...) {
throw;
}
}
void Plot::Private::retire_completed_frame(Render_Frame* frame) {
void Plot::Private::consume_completed_frame(Render_Frame* frame) {
if (!frame)
throw std::invalid_argument("Plot received a null completed frame");
std::size_t slot_index{scene_frame_capacity};
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
Managed_Frame* managed{};
for (auto& slot : frame_slots) {
auto* address = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
slot.frame);
if (address != frame) continue;
auto expected = Frame_State::in_flight;
if (!frame_slots[index].state.compare_exchange_strong(
expected, Frame_State::callback_retired,
std::memory_order_acq_rel, std::memory_order_acquire))
throw std::logic_error("completed Plot frame is not in flight");
slot_index = index;
if (slot.state.load(std::memory_order_acquire) != Frame_State::consuming)
throw std::logic_error("completed Plot frame was not published");
managed = &slot;
break;
}
if (slot_index == scene_frame_capacity)
throw std::logic_error(
"frame callback has no externally owned active frame");
auto no_retired = scene_frame_capacity;
if (!callback_retired_slot.compare_exchange_strong(
no_retired, slot_index, std::memory_order_acq_rel,
if (!managed)
throw std::logic_error("frame callback has no owned Plot frame");
/*
* Scene 已在调用本 callback 前释放自己的 render admission;这里同步
* 释放 Plot 的 view/update 门,并立刻唤醒 busy 期间保留的 latest tick。
* 之后外接 DAG 仍在当前 Frame trace 内执行,但不会阻塞下一帧渲染。
*/
release_render_admission(lifetime);
if (post_publish_graph.empty()) return;
bool expected = false;
if (!post_publish_busy.compare_exchange_strong(
expected, true, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("Plot has more than one callback-retired frame");
if (frame->taskflow_trace_requested()) {
auto trace = std::make_shared<const Taskflow_Frame_Trace>(
frame->taskflow_trace());
auto control = taskflow_trace_control.load(std::memory_order_acquire);
for (;;) {
const auto requested = static_cast<std::uint32_t>(control >> 32U);
const auto captured = static_cast<std::uint32_t>(control);
if (captured >= requested) break;
taskflow_trace_slots[captured].store(trace,
std::memory_order_release);
const auto next = (static_cast<std::uint64_t>(requested) << 32U) |
static_cast<std::uint64_t>(captured + 1U);
if (taskflow_trace_control.compare_exchange_weak(
control, next, std::memory_order_release,
std::memory_order_acquire))
break;
return;
/*
* 外接图异步提交;上一轮尚未完成时直接合并到 sampler 内的 latest,绝不
* 在 Taskflow Worker 内等待。诊断使用独立 Render_Frame 保存外接图观察
* 窗口,因此 Scene 的物理帧可立即退役并被下一次渲染复用。
*/
const bool local_post_publish_trace = mark_post_publish_taskflow_trace();
auto trace_frame = local_post_publish_trace
? std::make_shared<Render_Frame>(frame->identity())
: std::shared_ptr<Render_Frame>{};
bool local_trace_started{};
if (trace_frame) {
trace_frame->request_taskflow_trace();
local_trace_started =
aethera::detail::begin_taskflow_trace(*trace_frame);
}
const auto weak = lifetime;
auto completion = [weak, trace_frame, local_trace_started] {
const auto owner = weak.lock();
if (!owner) return;
if (local_trace_started) {
aethera::detail::finish_taskflow_trace(*trace_frame);
owner->d->store_trace(
owner->d->post_publish_trace_control,
owner->d->post_publish_trace_slots,
trace_frame->taskflow_trace());
}
owner->d->post_publish_busy.store(false, std::memory_order_release);
};
try {
if (trace_frame)
aethera::detail::run_taskflow(
post_publish_graph, *trace_frame, "plot.post_publish",
std::move(completion));
else
aethera::detail::run_taskflow(
post_publish_graph, std::move(completion));
}
if (frame_policy.frame_completed()) {
const auto weak = lifetime;
aethera::schedule_task("web.plot.render.deferred", [weak] {
const auto owner = weak.lock();
if (!owner) return;
try {
const auto now = std::chrono::steady_clock::now();
const auto elapsed = now - owner->d->clock_origin;
owner->d->render_frame(Plot_Render_Tick{
now, 0,
std::chrono::duration<double, std::milli>(elapsed).count()});
}
catch (...) {
owner->d->fail(std::current_exception());
}
});
catch (...) {
if (local_trace_started)
aethera::detail::finish_taskflow_trace(*trace_frame);
if (local_post_publish_trace)
post_publish_trace_remaining.fetch_add(1, std::memory_order_release);
post_publish_busy.store(false, std::memory_order_release);
throw;
}
}
void Plot::Private::retire_completed_frame(Render_Frame* frame) {
if (!frame)
throw std::invalid_argument("Plot received a null retired frame");
Managed_Frame* managed{};
for (auto& slot : frame_slots) {
auto* address = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
slot.frame);
if (address != frame) continue;
managed = &slot;
break;
}
if (!managed)
throw std::logic_error("retired frame has no owned Plot slot");
{
std::lock_guard lock(completed_frame_statistics_mutex);
completed_frame_statistics_state = completed_frame_statistics.submit(
*frame, std::holds_alternative<std::unique_ptr<Frame_2D>>(managed->frame)
? Frame_Dimension::two_dimensional
: Frame_Dimension::three_dimensional);
}
if (frame->taskflow_trace_requested())
store_trace(taskflow_trace_control, taskflow_trace_slots,
frame->taskflow_trace());
auto expected = Frame_State::consuming;
if (!managed->state.compare_exchange_strong(
expected, Frame_State::available, std::memory_order_acq_rel,
std::memory_order_acquire))
throw std::logic_error("retired Plot frame is not consuming");
/* 三个消费者槽曾全部占满时,退役一个槽后继续 latest pending。 */
arm_tick_consumer(lifetime);
}
void Plot::Private::attach_completion(
std::unique_ptr<Task_Graph> completion) {
if (!completion || completion->empty())
throw std::invalid_argument("Plot completion pipeline is empty");
auto& scene_completion = std::visit(
[](auto& scene_value) -> Task_Graph& {
return scene_value->completion_taskflow();
}, scene);
completion_extensions.push_back(std::move(completion));
auto extension = scene_completion.compose(
auto extension = post_publish_graph.compose(
completion_extensions.back()->name(), *completion_extensions.back());
extension.describe("owner", "scene")
.describe("stage", "frame pipeline extension");
completion_tail.precede(extension);
completion_tail = std::move(extension);
extension.describe("owner", "plot")
.describe("stage", "post-publish frame pipeline extension");
if (has_post_publish_tail) post_publish_tail.precede(extension);
post_publish_tail = std::move(extension);
has_post_publish_tail = true;
}
Plot::Plot(std::unique_ptr<Scene_2D> scene,
@@ -838,29 +1019,38 @@ void Plot::ensure_started() {
tick.issued_at, tick.sequence, tick.time_milliseconds});
}
});
d->refresh_schedule();
/*
* Kernel Frame_Scheduler 只产生每个 Scene 独立的周期 tick;
* Scene::render(Frame*) 与外部 Frame 所有权保持原样。完成回调先归还
* 当前帧;若期间出现 immediate 请求,则回调返回后重新投递 Taskflow
* 2D 的 callback 在 Scene render admission 已释放后运行:先执行所有
* post-publish 外接 DAGScene 完成 frame_ready 与 trace 收口后,再由
* retired callback 归还物理槽。这样 H264(N) 可与 Render(N+1) 重叠
*/
if (auto* scene = std::get_if<std::unique_ptr<Scene_2D>>(&d->scene)) {
(*scene)->set_frame_callback(
[weak](Frame_2D* frame) {
if (auto owner = weak.lock()) {
try { owner->d->retire_completed_frame(frame); }
catch (...) { owner->d->fail(std::current_exception()); }
}
});
(*scene)->set_frame_callback([weak](Frame_2D* frame) {
if (auto owner = weak.lock()) {
try { owner->d->consume_completed_frame(frame); }
catch (...) { owner->d->fail(std::current_exception()); }
}
});
(*scene)->set_frame_retired_callback([weak](Frame_2D* frame) {
if (auto owner = weak.lock()) {
try { owner->d->retire_completed_frame(frame); }
catch (...) { owner->d->fail(std::current_exception()); }
}
});
} else {
std::get<std::unique_ptr<Scene_3D>>(d->scene)->set_frame_callback(
[weak](Frame_3D* frame) {
if (auto owner = weak.lock()) {
try { owner->d->retire_completed_frame(frame); }
catch (...) { owner->d->fail(std::current_exception()); }
if (auto owner = weak.lock()) {
try {
owner->d->consume_completed_frame(frame);
owner->d->retire_completed_frame(frame);
}
catch (...) { owner->d->fail(std::current_exception()); }
}
});
}
/* callback 必须先于周期时钟安装,避免首帧在初始化窗口进入 Scene。 */
d->refresh_schedule();
});
}
@@ -935,41 +1125,17 @@ void Plot::schedule_render(Plot_Render_Tick tick) {
ensure_started();
if (d->terminal_failure.load(std::memory_order_acquire)) return;
d->received_tick_count.fetch_add(1, std::memory_order_relaxed);
const auto next = std::make_shared<const Plot_Render_Tick>(std::move(tick));
if (d->pending_tick.exchange(next, std::memory_order_acq_rel))
d->coalesced_tick_count.fetch_add(1, std::memory_order_relaxed);
if (d->tick_task_scheduled.exchange(true, std::memory_order_acq_rel)) return;
const auto weak = weak_from_this();
aethera::schedule_task("web.plot.tick.consume", [weak] {
const auto owner = weak.lock();
if (!owner) return;
try {
owner->d->consume_tick(weak);
}
catch (...) {
owner->d->fail(std::current_exception());
}
});
d->defer_tick(tick);
d->arm_tick_consumer(weak_from_this());
}
void Plot::render_once() {
ensure_started();
if (!d->frame_policy.request_immediate()) return;
const auto weak = weak_from_this();
aethera::schedule_task("web.plot.render.once", [weak] {
const auto owner = weak.lock();
if (!owner) return;
try {
const auto elapsed = std::chrono::steady_clock::now() -
owner->d->clock_origin;
owner->d->render_frame(Plot_Render_Tick{
std::chrono::steady_clock::now(),
0, std::chrono::duration<double, std::milli>(elapsed).count()});
}
catch (...) {
owner->d->fail(std::current_exception());
}
});
const auto now = std::chrono::steady_clock::now();
const auto elapsed = now - d->clock_origin;
schedule_render(Plot_Render_Tick{
now, 0, std::chrono::duration<double, std::milli>(elapsed).count()});
}
void Plot::submit_input(Plot_Input_Event event) {
@@ -1021,10 +1187,25 @@ nlohmann::json Plot::diagnostics() const {
std::uint64_t dropped_sequences{};
double frame_rate{};
bool is_3d{};
const auto read_statistics = [&](const auto& state) {
const auto& statistics = state.frame_statistics;
append_statistic_json(frame_statistics, statistics);
const auto read_scene_statistics = [&](const auto& state) {
append_event_statistics_json(input_statistics, state.event_statistics);
};
std::visit([&](const auto& scene) {
using Scene_Pointer = std::remove_cvref_t<decltype(scene)>;
if constexpr (std::same_as<Scene_Pointer, std::unique_ptr<Scene_2D>>) {
scene->template access_state<Render_Scene_2D::Base_Tag>(
read_scene_statistics);
} else {
is_3d = true;
scene->template access_state<Render_Scene_3D::Base_Tag>(
read_scene_statistics);
}
}, d->scene);
{
std::lock_guard lock(d->completed_frame_statistics_mutex);
const auto statistics = d->completed_frame_statistics_state;
append_statistic_json(frame_statistics, statistics);
identity = statistics.identity;
created_time_unix_ns = statistics.created_time_unix_ns;
dropped_sequences = statistics.dropped_sequences;
@@ -1032,18 +1213,7 @@ nlohmann::json Plot::diagnostics() const {
static_cast<std::size_t>(Frame_Statistic::frame_interval_ms)];
frame_rate = interval.trimmed_average > 0.0
? 1'000.0 / interval.trimmed_average : 0.0;
};
std::visit([&](const auto& scene) {
using Scene_Pointer = std::remove_cvref_t<decltype(scene)>;
if constexpr (std::same_as<Scene_Pointer, std::unique_ptr<Scene_2D>>) {
scene->template access_state<Render_Scene_2D::Base_Tag>(
read_statistics);
} else {
is_3d = true;
scene->template access_state<Render_Scene_3D::Base_Tag>(
read_statistics);
}
}, d->scene);
}
const auto pacing = d->frame_policy.read();
const auto stream = d->stream_snapshot();
@@ -1057,8 +1227,7 @@ nlohmann::json Plot::diagnostics() const {
}
const auto format = is_3d
? pixel_format_name(Frame_3D::native_pixel_format)
: pixel_format_name(pacing.video_enabled
? render_2d::Pixel_Format::rgba8 : Frame_2D::native_pixel_format);
: pixel_format_name(Frame_2D::native_pixel_format);
const auto native_format = is_3d
? pixel_format_name(Frame_3D::native_pixel_format)
: pixel_format_name(Frame_2D::native_pixel_format);
@@ -1090,6 +1259,8 @@ nlohmann::json Plot::diagnostics() const {
{"coalesced_ticks", d->coalesced_tick_count.load(std::memory_order_relaxed)},
{"policy_skips", d->policy_skip_count.load(std::memory_order_relaxed)},
{"preparation_busy", d->preparation_busy_count.load(std::memory_order_relaxed)},
{"deferred_resumes", d->deferred_resume_count.load(std::memory_order_relaxed)},
{"render_admission_busy", d->render_admission_busy.load(std::memory_order_relaxed)},
{"frame_slot_busy", d->frame_slot_busy_count.load(std::memory_order_relaxed)},
{"scene_rejections", d->scene_rejection_count.load(std::memory_order_relaxed)},
{"submitted_frames", d->submitted_frame_count.load(std::memory_order_relaxed)}}},
@@ -1149,26 +1320,45 @@ void Plot::request_taskflow_trace(std::size_t frame_count) {
}
nlohmann::json Plot::taskflow_trace() const {
nlohmann::json frames = nlohmann::json::array();
const auto control = d->taskflow_trace_control.load(
std::memory_order_acquire);
const auto requested = static_cast<std::uint32_t>(control >> 32U);
const auto captured = static_cast<std::uint32_t>(control);
for (std::uint32_t index = 0; index < captured; ++index)
if (const auto trace = d->taskflow_trace_slots[index].load(
return d->trace_response(d->taskflow_trace_control,
d->taskflow_trace_remaining,
d->taskflow_trace_slots);
}
void Plot::request_post_publish_taskflow_trace(std::size_t frame_count) {
if (frame_count == 0 ||
frame_count > Private::maximum_taskflow_trace_frames)
throw std::invalid_argument(
"Taskflow post-publish trace frame_count must be between 1 and 120");
ensure_started();
auto control = d->post_publish_trace_control.load(std::memory_order_acquire);
for (;;) {
const auto requested = static_cast<std::uint32_t>(control >> 32U);
const auto captured = static_cast<std::uint32_t>(control);
if (requested != captured)
throw std::logic_error(
"A post-publish Taskflow trace request is already active");
const auto next = static_cast<std::uint64_t>(frame_count) << 32U;
if (d->post_publish_trace_control.compare_exchange_weak(
control, next, std::memory_order_release,
std::memory_order_acquire))
frames.push_back(taskflow_trace_json(*trace));
const auto remaining = d->taskflow_trace_remaining.load(
std::memory_order_acquire);
return {
{"protocol", "aethera.taskflow.frames"}, {"version", 1},
{"requested", requested}, {"remaining", remaining},
{"captured", frames.size()},
{"complete", requested != 0 && frames.size() == requested},
{"frames", std::move(frames)}};
break;
}
for (auto& slot : d->post_publish_trace_slots)
slot.store({}, std::memory_order_release);
d->post_publish_trace_remaining.store(frame_count, std::memory_order_release);
}
nlohmann::json Plot::post_publish_taskflow_trace() const {
return d->trace_response(d->post_publish_trace_control,
d->post_publish_trace_remaining,
d->post_publish_trace_slots);
}
void Plot::reset_diagnostics() {
std::visit([](auto& scene) { scene->reset_frame_statistics(); }, d->scene);
std::lock_guard lock(d->completed_frame_statistics_mutex);
d->completed_frame_statistics.reset();
d->completed_frame_statistics_state = {};
}
}