@@ -19,11 +19,11 @@
# include <initializer_list>
# include <limits>
# include <memory>
# include <mutex>
# include <optional>
# include <span>
# include <stdexcept>
# include <unordered_map>
# include <unordered_set>
# include <utility>
# include <variant>
# include <vector>
@@ -200,7 +200,7 @@ void append_event_statistics_json(nlohmann::json& output,
nlohmann : : json taskflow_trace_json (
const Taskflow_Frame_Trace & trace ,
const nlohmann : : json & component_snapshots ) {
const nlohmann : : json & captured_components ) {
nlohmann : : json markers = nlohmann : : json : : object ( ) ;
for ( const auto & marker : trace . markers )
markers [ magic_enum : : enum_name ( marker . marker ) ] =
@@ -228,14 +228,14 @@ nlohmann::json taskflow_trace_json(
{ " attributes " , std : : move ( attributes ) } } ;
const auto owner = encoded [ " attributes " ] . value (
" owner_component " , std : : string { } ) ;
if ( ! owner . empty ( ) & & component_snapshots . contains ( owner ) ) {
const auto & snapshot = component_snapshots . at ( owner ) ;
if ( ! owner . empty ( ) & & captured_components . contains ( owner ) ) {
const auto & captured = captured_components . at ( owner ) ;
encoded [ " owner " ] = {
{ " component " , owner } ,
{ " label " , snapshot . value ( " label " , owner ) } ,
{ " kind " , snapshot . value ( " kind " , std : : string { } ) } } ;
encoded [ " prop " ] = snapshot . value ( " prop " , nlohmann : : json : : object ( ) ) ;
encoded [ " state " ] = snapshot . value ( " state " , nlohmann : : json : : object ( ) ) ;
{ " label " , captured . value ( " label " , owner ) } ,
{ " kind " , captured . value ( " kind " , std : : string { } ) } } ;
encoded [ " prop " ] = captured . value ( " prop " , nlohmann : : json : : object ( ) ) ;
encoded [ " state " ] = captured . value ( " state " , nlohmann : : json : : object ( ) ) ;
}
nodes . push_back ( std : : move ( encoded ) ) ;
}
@@ -343,6 +343,10 @@ struct Plot::Private {
std : : chrono : : microseconds presentation_time { } ; /* 共享页面时钟产生的媒体时间戳。 */
Frame frame { } ; /* 三缓冲物理槽拥有且反复承载逻辑帧。 */
std : : atomic < Frame_State > state { Frame_State : : available } ; /* 本槽唯一生命周期状态。 */
std : : atomic_size_t diagnostic_readers { } ; /* 无锁诊断读取认领;非零时该槽不可复用。 */
std : : uint64_t statistics_generation { } ; /* 与本槽中完成帧统计共同发布。 */
Frame_Statistics_State statistics { } ; /* 统计结果归属当前完成帧,不在 Plot 复制。 */
std : : atomic < Managed_Frame * > retired_next { } ; /* 无锁退役队列的槽内侵入链接。 */
} ;
struct Consumer {
@@ -370,6 +374,9 @@ struct Plot::Private {
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 < Managed_Frame * > latest_statistics_frame { } ; /* 只定位权威帧槽,不保存统计副本。 */
std : : atomic < Managed_Frame * > retired_frames { } ; /* 多回调生产、唯一 Taskflow 任务消费。 */
std : : atomic_bool retired_frame_task_scheduled { } ;
Scene scene ; /* 析构顺序保证 Scene 先停止,再释放物理帧。 */
std : : atomic < std : : shared_ptr < const Plot_Render_Tick > > pending_tick { } ;
std : : atomic_bool tick_task_scheduled { } ; /* 唯一短任务准入;不占用 Worker 等待。 */
@@ -401,8 +408,8 @@ struct Plot::Private {
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 { } ;
std : : uint64_t applied_statistics_generation { } ; /* 仅完成帧退役任务读写。 */
std : : atomic_uint64_t statistics_generation { } ; /* reset 只推进代次,不触碰单写者累加器。 */
template < typename Scene_Object >
Private ( std : : unique_ptr < Scene_Object > value_scene ,
@@ -438,6 +445,9 @@ struct Plot::Private {
void publish_completed_frame ( ) ;
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 ( ) ;
@@ -445,7 +455,7 @@ struct Plot::Private {
std : : array < std : : atomic < std : : shared_ptr < const nlohmann : : json > > ,
maximum_taskflow_trace_frames > & slots ,
const Taskflow_Frame_Trace & trace ,
const nlohmann : : json & component_snapshots = { } ) ;
const nlohmann : : json & captured_components = { } ) ;
[[nodiscard]] nlohmann : : json trace_response (
const std : : atomic_uint64_t & control ,
const std : : atomic_size_t & remaining ,
@@ -628,9 +638,9 @@ void Plot::Private::store_trace(
std : : array < std : : atomic < std : : shared_ptr < const nlohmann : : json > > ,
maximum_taskflow_trace_frames > & slots ,
const Taskflow_Frame_Trace & value ,
const nlohmann : : json & component_snapshots ) {
const nlohmann : : json & captured_components ) {
auto trace = std : : make_shared < const nlohmann : : json > (
taskflow_trace_json ( value , component_snapshots ) ) ;
taskflow_trace_json ( value , captured_components ) ) ;
auto state = control . load ( std : : memory_order_acquire ) ;
for ( ; ; ) {
const auto requested = static_cast < std : : uint32_t > ( state > > 32U ) ;
@@ -690,11 +700,32 @@ void Plot::Private::render_frame(Plot_Render_Tick tick) {
* 诊 断 生 命 周 期 。 只 要 还 有 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 : : rendering ,
std : : memory_order_acq_rel , std : : memory_order_acquire ) )
auto * candidate = & frame_slots [ index ] ;
auto * published = candidate ;
const bool was_latest = latest_statistics_frame . compare_exchange_strong (
published , nullptr , std : : memory_order_acq_rel ,
std : : memory_order_acquire ) ;
if ( candidate - > diagnostic_readers . load ( std : : memory_order_acquire ) ! = 0 ) {
if ( was_latest ) {
Managed_Frame * empty { } ;
static_cast < void > ( latest_statistics_frame . compare_exchange_strong (
empty , candidate , std : : memory_order_release ,
std : : memory_order_relaxed ) ) ;
}
continue ;
}
auto expected = Frame_State : : available ;
if ( ! candidate - > state . compare_exchange_strong (
expected , Frame_State : : rendering ,
std : : memory_order_acq_rel , std : : memory_order_acquire ) ) {
if ( was_latest ) {
Managed_Frame * empty { } ;
static_cast < void > ( latest_statistics_frame . compare_exchange_strong (
empty , candidate , std : : memory_order_release ,
std : : memory_order_relaxed ) ) ;
}
continue ;
}
slot_index = index ;
managed = & frame_slots [ index ] ;
break ;
@@ -961,6 +992,46 @@ 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 ) ;
if ( ! list ) break ;
std : : vector < Managed_Frame * > frames ;
while ( list ) {
auto * next = list - > retired_next . exchange (
nullptr , std : : memory_order_relaxed ) ;
frames . push_back ( list ) ;
list = next ;
}
std : : ranges : : sort ( frames , { } , [ ] ( const Managed_Frame * managed ) {
return std : : visit (
[ ] ( const auto & value ) { return value - > identity ( ) . sequence ; } ,
managed - > frame ) ;
} ) ;
for ( auto * managed : frames ) {
auto * frame = std : : visit (
[ ] ( const auto & value ) - > Render_Frame * { return value . get ( ) ; } ,
managed - > frame ) ;
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 ) {
if ( ! frame )
throw std : : invalid_argument ( " Plot received a null retired frame " ) ;
@@ -976,23 +1047,82 @@ void Plot::Private::retire_completed_frame(Render_Frame* frame) {
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 ) ;
}
auto * head = retired_frames . load ( std : : memory_order_relaxed ) ;
do {
managed - > retired_next . store ( head , std : : memory_order_relaxed ) ;
} while ( ! retired_frames . compare_exchange_weak (
head , managed , std : : memory_order_release ,
std : : memory_order_relaxed ) ) ;
arm_retired_frame_consumer ( ) ;
}
if ( frame - > taskflow_trace_requested ( ) )
store_trace ( taskflow_trace_control , taskflow_trace_slots ,
frame - > take_taskflow_trace ( ) , view - > component_snapshots ( ) ) ;
void Plot : : Private : : finalize_retired_frame ( Render_Frame * 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 ) {
managed = & slot ;
break ;
}
}
if ( ! managed )
throw std : : logic_error ( " retired frame has no owned Plot slot " ) ;
const auto statistics_generation_value =
statistics_generation . load ( std : : memory_order_acquire ) ;
if ( applied_statistics_generation ! = statistics_generation_value ) {
completed_frame_statistics . reset ( ) ;
applied_statistics_generation = statistics_generation_value ;
}
managed - > statistics = completed_frame_statistics . submit ( * frame ) ;
managed - > statistics_generation = statistics_generation_value ;
std : : optional < Taskflow_Frame_Trace > captured_trace ;
nlohmann : : json captured_components ;
if ( frame - > taskflow_trace_requested ( ) ) {
captured_trace . emplace ( frame - > take_taskflow_trace ( ) ) ;
std : : unordered_set < std : : uint64_t > executed_nodes ;
for ( const auto & execution : captured_trace - > tasks )
executed_nodes . insert ( execution . native_id ) ;
std : : vector < std : : string > executed_components ;
for ( const auto & graph : captured_trace - > graphs ) {
for ( const auto & node : graph . nodes ) {
if ( ! executed_nodes . contains ( node . native_id ) ) continue ;
const auto owner = std : : ranges : : find (
node . attributes , " owner_component " ,
& std : : pair < std : : string , std : : string > : : first ) ;
if ( owner = = node . attributes . end ( ) | | owner - > second . empty ( ) | |
std : : ranges : : find ( executed_components , owner - > second ) ! =
executed_components . end ( ) ) continue ;
executed_components . push_back ( owner - > second ) ;
}
}
captured_components = view - > capture_components ( executed_components ) ;
}
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_statistics_frame . store ( managed , std : : memory_order_release ) ;
/* 三个消费者槽曾全部占满时,退役一个槽后继续 latest pending。 */
arm_tick_consumer ( lifetime ) ;
if ( captured_trace ) {
auto owner = lifetime ;
schedule_task ( " plot.taskflow.serialize " , [ owner ,
trace = std : : move ( * captured_trace ) ,
components = std : : move ( captured_components ) ] ( ) mutable {
if ( const auto plot = owner . lock ( ) )
plot - > d - > store_trace (
plot - > d - > taskflow_trace_control ,
plot - > d - > taskflow_trace_slots , trace , components ) ;
} ) ;
}
}
void Plot : : Private : : attach_completion (
@@ -1221,8 +1351,23 @@ nlohmann::json Plot::diagnostics() const {
} , d - > scene ) ;
{
std : : lock_guard lock ( d - > completed_frame_statistics_mutex ) ;
const auto statistics = d - > completed_frame_statistics_state ;
const auto generation =
d - > statistics_generation . load ( std : : memory_order_acquire ) ;
auto * completed_frame =
d - > latest_statistics_frame . load ( std : : memory_order_acquire ) ;
Frame_Statistics_State statistics { } ;
if ( completed_frame ) {
completed_frame - > diagnostic_readers . fetch_add (
1 , std : : memory_order_acq_rel ) ;
if ( d - > latest_statistics_frame . load ( std : : memory_order_acquire ) = =
completed_frame & &
completed_frame - > state . load ( std : : memory_order_acquire ) = =
Private : : Frame_State : : available & &
completed_frame - > statistics_generation = = generation )
statistics = completed_frame - > statistics ;
completed_frame - > diagnostic_readers . fetch_sub (
1 , std : : memory_order_release ) ;
}
append_statistic_json ( frame_statistics , statistics ) ;
identity = statistics . identity ;
created_time_unix_ns = statistics . created_time_unix_ns ;
@@ -1375,8 +1520,6 @@ nlohmann::json Plot::post_publish_taskflow_trace() const {
void Plot : : reset_diagnostics ( ) {
std : : visit ( [ ] ( auto & scene ) { scene - > reset_diagnostics ( ) ; } , d - > scene ) ;
std : : lock_guard lock ( d - > completed_frame_statistics_mutex ) ;
d - > completed_frame_statistics . reset ( ) ;
d - > completed_frame_statistics_state = { } ;
d - > statistics_generation . fetch_add ( 1 , std : : memory_order_acq_rel ) ;
}
}