加所有权之前

This commit is contained in:
2026-08-28 12:14:15 +08:00
parent 81d3d50966
commit 0bc281b6b9
86 changed files with 12167 additions and 427 deletions
+123 -92
View File
@@ -298,6 +298,9 @@ nlohmann::json datoviz_observation_json(
{"path", magic_enum::enum_name(value.path)},
{"gpu_timing_requested", value.gpu_timing_requested},
{"readback_requested", value.readback_requested},
{"controller_input_applied", value.controller_input_applied},
{"prepare_released_after_submission",
value.prepare_released_after_submission},
{"timings_ms", {
{"render_domain_queue_wait",
milliseconds(value.render_domain_queue_wait_ns)},
@@ -421,8 +424,24 @@ nlohmann::json taskflow_trace_json(
namespace {
[[nodiscard]] Event_Timeline_Time input_timeline_time(
double time_milliseconds) {
constexpr long double nanoseconds_per_millisecond{1'000'000.0L};
constexpr long double maximum_milliseconds =
static_cast<long double>(std::numeric_limits<std::uint64_t>::max()) /
nanoseconds_per_millisecond;
if (!std::isfinite(time_milliseconds) || time_milliseconds < 0.0 ||
static_cast<long double>(time_milliseconds) > maximum_milliseconds)
throw std::invalid_argument(
"input time_milliseconds must be finite, non-negative and representable");
return Event_Timeline_Time{static_cast<std::uint64_t>(
static_cast<long double>(time_milliseconds) *
nanoseconds_per_millisecond)};
}
template <typename Scene_Object>
void dispatch_plot_input(Scene_Object& scene, const Plot_Input_Event& input) {
const auto occurred_at = input_timeline_time(input.time_milliseconds);
const auto dispatch = [&](auto event) {
scene.template submit_stream<aethera::Scene_Event_Stream_Tag>(std::move(event));
};
@@ -438,13 +457,14 @@ void dispatch_plot_input(Scene_Object& scene, const Plot_Input_Event& input) {
case Event_Type::pointer_press:
case Event_Type::pointer_release: {
auto event = scene.template make_event<Basic_Pointer_Event<Point_F>>(
input.type);
input.type, occurred_at);
apply_pointer(*event);
dispatch(std::move(event));
break;
}
case Event_Type::wheel: {
auto event = scene.template make_event<Basic_Wheel_Event<Point_F>>();
auto event = scene.template make_event<Basic_Wheel_Event<Point_F>>(
occurred_at);
apply_pointer(*event);
event->pixel_delta_x = input.pixel_delta_x;
event->pixel_delta_y = input.pixel_delta_y;
@@ -455,7 +475,8 @@ void dispatch_plot_input(Scene_Object& scene, const Plot_Input_Event& input) {
}
case Event_Type::key_press:
case Event_Type::key_release: {
auto event = scene.template make_event<Key_Event>(input.type);
auto event = scene.template make_event<Key_Event>(
input.type, occurred_at);
event->key = input.key;
event->native_key = input.native_key;
event->modifiers = input.modifiers;
@@ -464,7 +485,7 @@ void dispatch_plot_input(Scene_Object& scene, const Plot_Input_Event& input) {
break;
}
default:
dispatch(scene.template make_event<Event>(input.type));
dispatch(scene.template make_event<Event>(input.type, occurred_at));
break;
}
}
@@ -523,10 +544,9 @@ struct Plot::Private {
std::unique_ptr<Frame_Policy_3D> frame_policy_3d{};
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::array<Managed_Frame, scene_frame_capacity> frame_slots{}; /* 帧策略拥有并反复调度的稳定三缓冲;Scene 只借用。 */
std::atomic<Managed_Frame*> latest_statistics_frame{}; /* 只定位权威帧槽,不保存统计副本。 */
std::atomic<Managed_Frame*> retired_frames{}; /* 多回调生产、唯一 Taskflow 任务消费。 */
std::atomic_bool retired_frame_task_scheduled{};
std::atomic<Managed_Frame*> retired_frames{}; /* 完成回调返回、唯一帧策略写者消费的物理帧。 */
Scene scene; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */
moodycamel::ConcurrentQueue<Plot_Render_Tick> tick_requests{}; /* 多生产者提交、唯一短任务消费的帧请求流。 */
std::optional<Plot_Render_Tick> deferred_tick{}; /* 仅 tick consumer 任务访问的 latest 延后请求。 */
@@ -545,7 +565,6 @@ struct Plot::Private {
std::atomic_uint64_t post_publish_trace_control{}; /* 高 32 位 requested,低 32 位 captured。 */
std::array<std::atomic<std::shared_ptr<const nlohmann::json>>,
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{};
@@ -580,15 +599,6 @@ struct Plot::Private {
else
slot.frame = std::make_unique<Frame_3D>(Frame_Identity{});
}
if constexpr (std::same_as<Scene_Object, Scene_2D>) {
auto& completion = std::get<std::unique_ptr<Scene_2D>>(scene)
->completion_taskflow();
completion_tail = completion.add("plot.frame.publish", [this] {
publish_completed_frame();
});
completion_tail.describe("owner", "plot")
.describe("stage", "completed pixels publish");
}
}
template <typename Callback>
@@ -612,11 +622,8 @@ struct Plot::Private {
}
[[nodiscard]] bool consume_frame_policy_events() {
const bool configuration_changed = with_frame_policy(
return with_frame_policy(
[](auto& policy) { return policy.consume_events(); });
if (frame_policy_3d)
static_cast<void>(frame_policy_3d->consume_3d_events());
return configuration_changed;
}
[[nodiscard]] nlohmann::json schema() const;
@@ -630,12 +637,11 @@ struct Plot::Private {
void refresh_schedule();
void clock_tick(const Plot_Render_Tick& tick);
void render_frame(Plot_Render_Tick tick);
void publish_completed_frame(Render_Frame* completed = nullptr);
void publish_completed_frame(Render_Frame* completed);
void consume_completed_frame(Render_Frame* frame);
void retire_completed_frame(Render_Frame* frame);
void consume_retired_frames();
void finalize_retired_frame(Render_Frame* frame);
void arm_retired_frame_consumer();
void attach_completion(std::unique_ptr<Task_Graph> completion);
[[nodiscard]] bool mark_taskflow_trace(Render_Frame& frame);
[[nodiscard]] bool mark_post_publish_taskflow_trace();
@@ -826,6 +832,7 @@ void Plot::Private::release_render_admission(std::weak_ptr<Plot> lifetime) {
}
void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
consume_retired_frames();
if (consume_frame_policy_events()) refresh_schedule();
const auto observed_generation =
tick_request_generation.load(std::memory_order_acquire);
@@ -837,10 +844,12 @@ void Plot::Private::consume_tick(std::weak_ptr<Plot> lifetime) {
auto tick = std::exchange(deferred_tick, {});
if (tick) clock_tick(*tick);
}
consume_retired_frames();
if (consume_frame_policy_events()) refresh_schedule();
tick_task_scheduled.store(false, std::memory_order_release);
if (tick_request_generation.load(std::memory_order_acquire) !=
observed_generation ||
retired_frames.load(std::memory_order_acquire) ||
(render_admission.load(std::memory_order_acquire) ==
Render_Admission_State::ready && deferred_tick))
arm_tick_consumer(std::move(lifetime));
@@ -940,7 +949,12 @@ 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 = pacing_state();
if (!pacing.render_enabled || streams.consumers->empty()) return;
const bool direct_diagnostics_frame =
streams.consumers->empty() &&
tick.source == Frame_Request_Source::immediate;
if (!pacing.render_enabled ||
(streams.consumers->empty() && !direct_diagnostics_frame))
return;
auto admission_expected = Render_Admission_State::ready;
if (!render_admission.compare_exchange_strong(
@@ -1024,8 +1038,14 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
};
try {
tick.width = streams.width;
tick.height = streams.height;
if (streams.consumers->empty()) {
tick.width = std::clamp(tick.width, 160U, 1920U) & ~1U;
tick.height = std::clamp(tick.height, 120U, 1080U) & ~1U;
}
else {
tick.width = streams.width;
tick.height = streams.height;
}
const std::uint64_t sequence = next_frame_sequence++;
const Frame_Identity identity{
sequence, tick.sequence == 0 ? sequence : tick.sequence};
@@ -1128,23 +1148,23 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
void Plot::Private::publish_completed_frame(Render_Frame* completed) {
Render_Frame* frame{};
if (!completed)
throw std::invalid_argument("Plot received a null completed frame");
Managed_Frame* managed{};
for (std::size_t index = 0; index < frame_slots.size(); ++index) {
if (frame_slots[index].state.load(std::memory_order_acquire) !=
Frame_State::rendering)
continue;
if (managed)
throw std::logic_error("Plot has multiple frames in Scene rendering");
frame = std::visit(
for (auto& slot : frame_slots) {
auto* frame = std::visit(
[](const auto& value) -> Render_Frame* { return value.get(); },
frame_slots[index].frame);
if (completed && frame != completed) continue;
managed = &frame_slots[index];
if (completed) break;
slot.frame);
if (frame != completed) continue;
managed = &slot;
break;
}
if (!managed)
throw std::logic_error("Scene completion graph has no rendering Plot frame");
throw std::logic_error("completed frame has no owning Plot policy slot");
if (managed->state.load(std::memory_order_acquire) !=
Frame_State::rendering)
throw std::logic_error("completed Plot frame is not rendering");
auto* frame = completed;
const auto& pacing = pacing_state();
const auto identity = frame->identity();
@@ -1276,18 +1296,6 @@ void Plot::Private::consume_completed_frame(Render_Frame* frame) {
}
}
void Plot::Private::arm_retired_frame_consumer() {
bool expected{};
if (!retired_frame_task_scheduled.compare_exchange_strong(
expected, true, std::memory_order_acq_rel,
std::memory_order_acquire)) return;
const auto weak = lifetime;
schedule_task("web.plot.frame.retire", [weak] {
if (const auto owner = weak.lock())
owner->d->consume_retired_frames();
});
}
void Plot::Private::consume_retired_frames() {
for (;;) {
auto* list = retired_frames.exchange(nullptr, std::memory_order_acq_rel);
@@ -1311,9 +1319,6 @@ void Plot::Private::consume_retired_frames() {
finalize_retired_frame(frame);
}
}
retired_frame_task_scheduled.store(false, std::memory_order_release);
if (retired_frames.load(std::memory_order_acquire))
arm_retired_frame_consumer();
}
void Plot::Private::retire_completed_frame(Render_Frame* frame) {
@@ -1337,7 +1342,7 @@ void Plot::Private::retire_completed_frame(Render_Frame* frame) {
} while (!retired_frames.compare_exchange_weak(
head, managed, std::memory_order_release,
std::memory_order_relaxed));
arm_retired_frame_consumer();
arm_tick_consumer(lifetime);
}
void Plot::Private::finalize_retired_frame(Render_Frame* frame) {
@@ -1354,13 +1359,24 @@ void Plot::Private::finalize_retired_frame(Render_Frame* frame) {
if (!managed)
throw std::logic_error("retired frame has no owned Plot slot");
with_frame_policy([&](auto& policy) {
policy.record_frame_completed(
frame->identity().sequence,
managed->submitted_at.time_since_epoch().count() == 0
? std::chrono::steady_clock::duration::zero()
: std::chrono::steady_clock::now() - managed->submitted_at);
});
const auto completion_latency =
managed->submitted_at.time_since_epoch().count() == 0
? std::chrono::steady_clock::duration::zero()
: std::chrono::steady_clock::now() - managed->submitted_at;
std::optional<Datoviz_Frame_Observation> datoviz_observation;
if (auto* frame_3d = dynamic_cast<Frame_3D*>(frame)) {
datoviz_observation = frame_3d->take_datoviz_observation();
if (!datoviz_observation)
throw std::logic_error(
"completed 3D frame has no Datoviz observation");
frame_policy_3d->record_frame_completed(
frame->identity().sequence, completion_latency,
datoviz_observation->prepare_released_after_submission);
}
else {
frame_policy_2d->record_frame_completed(
frame->identity().sequence, completion_latency);
}
const auto statistics_generation_value =
statistics_generation.load(std::memory_order_acquire);
@@ -1393,10 +1409,8 @@ void Plot::Private::finalize_retired_frame(Render_Frame* frame) {
}
}
captured_components = view->capture_components(executed_components);
if (auto* frame_3d = dynamic_cast<render_3d::Frame_3D*>(frame)) {
if (auto observation = frame_3d->take_datoviz_observation())
captured_backend = datoviz_observation_json(*observation);
}
if (datoviz_observation)
captured_backend = datoviz_observation_json(*datoviz_observation);
}
auto expected = Frame_State::consuming;
@@ -1469,32 +1483,24 @@ void Plot::ensure_started() {
.source = Frame_Request_Source::periodic});
}
});
/*
* 2D 的 callback 在 Scene render admission 已释放后运行:先执行所有
* post-publish 外接 DAGScene 完成 frame_ready 与 trace 收口后,再由
* retired callback 归还物理槽。这样 H264(N) 可与 Render(N+1) 重叠。
*/
/* Scene 回调只表示借用结束;Frame 的发布、统计和复用始终由本帧策略决定。 */
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->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); }
try {
owner->d->publish_completed_frame(frame);
owner->d->consume_completed_frame(frame);
owner->d->retire_completed_frame(frame);
}
catch (...) { owner->d->fail(std::current_exception()); }
}
});
} else {
auto& scene_3d = std::get<std::unique_ptr<Scene_3D>>(d->scene);
scene_3d->set_submitted_frame_callback(
[weak](Frame_3D* frame, bool overlaps_gpu) {
[weak](Frame_3D*, bool) {
if (auto owner = weak.lock()) {
try {
owner->d->frame_policy_3d->record_gpu_submitted(
frame->identity().sequence, overlaps_gpu);
owner->d->release_render_admission(weak);
}
catch (...) { owner->d->fail(std::current_exception()); }
@@ -1505,8 +1511,6 @@ void Plot::ensure_started() {
if (auto owner = weak.lock()) {
try {
owner->d->publish_completed_frame(frame);
owner->d->frame_policy_3d->record_gpu_completed(
frame->identity().sequence);
owner->d->consume_completed_frame(frame);
owner->d->retire_completed_frame(frame);
}
@@ -1587,6 +1591,12 @@ void Plot::configure_stream(Stream_Id stream, std::uint32_t width,
}
void Plot::schedule_render(Plot_Render_Tick tick) {
if (!std::isfinite(tick.time_milliseconds) ||
tick.time_milliseconds < 0.0)
throw std::invalid_argument(
"render time_milliseconds must be finite and non-negative");
if (tick.width == 0 || tick.height == 0)
throw std::invalid_argument("render viewport must be non-zero");
ensure_started();
if (d->terminal_failure.load(std::memory_order_acquire)) return;
d->submit_tick_request(std::move(tick));
@@ -1605,6 +1615,7 @@ void Plot::render_once() {
}
void Plot::submit_input(Plot_Input_Event event) {
static_cast<void>(input_timeline_time(event.time_milliseconds));
ensure_started();
if (d->terminal_failure.load(std::memory_order_acquire)) return;
/*
@@ -1658,6 +1669,8 @@ nlohmann::json Plot::diagnostics() const {
Frame_Identity identity{};
std::uint64_t created_time_unix_ns{};
std::uint64_t dropped_sequences{};
std::uint32_t completed_width{};
std::uint32_t completed_height{};
double frame_rate{};
bool is_3d{};
const auto read_scene_statistics = [&](const auto& state) {
@@ -1688,8 +1701,27 @@ nlohmann::json Plot::diagnostics() const {
completed_frame &&
completed_frame->state.load(std::memory_order_acquire) ==
Private::Frame_State::available &&
completed_frame->statistics_generation == generation)
completed_frame->statistics_generation == generation) {
statistics = completed_frame->statistics;
std::visit([&](const auto& frame) {
using Frame_Pointer =
std::remove_cvref_t<decltype(frame)>;
if constexpr (std::same_as<
Frame_Pointer,
std::unique_ptr<Frame_2D>>) {
const auto image = frame->image();
completed_width = static_cast<std::uint32_t>(
std::max(0, image.width));
completed_height = static_cast<std::uint32_t>(
std::max(0, image.height));
}
else {
const auto extent = frame->extent();
completed_width = extent.width;
completed_height = extent.height;
}
}, completed_frame->frame);
}
completed_frame->diagnostic_readers.fetch_sub(
1, std::memory_order_release);
}
@@ -1729,8 +1761,12 @@ nlohmann::json Plot::diagnostics() const {
const auto native_format = is_3d
? pixel_format_name(Frame_3D::native_pixel_format)
: pixel_format_name(Frame_2D::native_pixel_format);
const auto pixel_width = completed_width == 0
? stream.width : completed_width;
const auto pixel_height = completed_height == 0
? stream.height : completed_height;
const std::size_t byte_length = pacing.video_enabled
? static_cast<std::size_t>(stream.width) * stream.height * 4U : 0U;
? static_cast<std::size_t>(pixel_width) * pixel_height * 4U : 0U;
nlohmann::json output{
{"protocol", "aethera.plot.diagnostics"}, {"version", 4},
{"dimension", is_3d ? "3D" : "2D"},
@@ -1744,7 +1780,7 @@ nlohmann::json Plot::diagnostics() const {
{"frame_rate_fps", frame_rate},
{"dropped_sequence_count", dropped_sequences},
{"window_capacity", diagnostic_window_capacity},
{"pixel", {{"width", stream.width}, {"height", stream.height},
{"pixel", {{"width", pixel_width}, {"height", pixel_height},
{"format", format}, {"native_format", native_format},
{"supported_formats", std::move(supported_formats)},
{"byte_length", byte_length}}},
@@ -1756,15 +1792,11 @@ nlohmann::json Plot::diagnostics() const {
const auto& pipeline = d->frame_policy_3d->read_state<
Frame_Policy_3D::Base_Tag>();
output["frame_policy"]["three_dimensional_pipeline"] = {
{"capacity", pipeline.pipeline_capacity},
{"gpu_submission_count", pipeline.gpu_submission_count},
{"capacity", Frame_Policy_3D::pipeline_capacity},
{"overlapped_release_count", pipeline.overlapped_release_count},
{"completion_gated_release_count",
pipeline.completion_gated_release_count},
{"gpu_completion_count", pipeline.gpu_completion_count},
{"gpu_in_flight", pipeline.gpu_in_flight},
{"peak_gpu_in_flight", pipeline.peak_gpu_in_flight},
{"last_submitted_sequence", pipeline.last_submitted_sequence},
{"last_completed_sequence", pipeline.last_completed_sequence}
};
const auto gpu = gpu_completion_state();
@@ -1859,7 +1891,6 @@ void Plot::reset_diagnostics() {
std::visit([](auto& scene) { scene->reset_diagnostics(); }, d->scene);
d->statistics_generation.fetch_add(1, std::memory_order_acq_rel);
d->with_frame_policy([](auto& policy) { policy.reset_statistics(); });
if (d->frame_policy_3d) d->frame_policy_3d->reset_3d_statistics();
d->arm_tick_consumer(weak_from_this());
}
}