diff --git a/docs/tensormap-and-ringbuffer-a2a3-vs-a5.md b/docs/tensormap-and-ringbuffer-a2a3-vs-a5.md index 53012416c9..846b3fd971 100644 --- a/docs/tensormap-and-ringbuffer-a2a3-vs-a5.md +++ b/docs/tensormap-and-ringbuffer-a2a3-vs-a5.md @@ -112,6 +112,7 @@ The functional differences group into the following themes: | URMA completion | A5-specific implementation and product capability gate | Yes, for now | Retain the A5 path; do not claim that URMA is available in the default build | | Next-block prefetch | A2/A3-only performance optimization | No | Retain on A2/A3; validate on A5 before considering a port | | Scheduler progress publication | AICPU topology and measured publication cost | No | Retain A5's 16-task batching; keep per-advance publication on A2/A3, where the portable implementation showed no significant benefit | +| Terminal deferred release | Measured end-of-run Scheduler release cost | No | Keep the current experiment scoped to A5; A2/A3 traces do not show the same large terminal release stall | | Fatal teardown | Software reliability strategy | No | Retain the current implementations; decide whether to converge after measuring the worst-case A5 teardown time | | Scheduler trace attribution | Software diagnostic strategy | No | Preserve the current traces; converge only after comparing generated timelines | @@ -293,6 +294,36 @@ measurements instead showed lower Effective time in all eight workloads, with an unweighted mean reduction of `2.81%`. Full A2/A3 measurements are recorded in the [PR benchmark follow-up](https://github.com/hw-native-sys/simpler/pull/1575#issuecomment-5310909143). +### Terminal Deferred Release: A5 Experiment Scope + +The terminal deferred-release experiment is currently A5-only. A5 swimlanes +showed Scheduler time extending beyond useful work while a large accumulated +release backlog updated task reference counts and advanced ring reclamation at +the end of a run. The corresponding A2/A3 traces reviewed for this work did +not show a comparable block of terminal release work, so there is no measured +A2/A3 bottleneck for this optimization to address. + +On A5, the experiment keeps exact per-task release while orchestration can +still create work. Once orchestration is sealed, Schedulers may discard their +local deferred-release backlogs at existing release boundaries. After every +task has completed and all Schedulers leave dispatch, the last thread at a +terminal barrier closes the remaining live ring slots in one pass and +publishes the final ring state. Errors and non-terminal exits retain exact +per-task lifecycle handling. + +The A5 benchmark covers seven workloads for 100 rounds and Qwen3 for five +rounds. No workload regressed by 5% in Orchestrator time. Batch Paged Attention, +the workload that exposed the terminal release cost, changed from +`5918.916 us` on the refreshed A5 Main baseline to `5819.991 us` with the +experiment (`-1.67%`). The same experiment had measured `6947.106 us` before +restoring the Main Scheduler entry gate, which also demonstrates that hot-code +layout must be preserved when evaluating the lifecycle optimization. + +These results establish an A5 optimization target, not a platform-independent +policy. A2/A3 keeps its existing release behavior unless a future A2/A3 +swimlane shows the same terminal release bottleneck and a separate benchmark +demonstrates a benefit. + ### Fatal Teardown The A2/A3 scheduler uses a dedicated fatal latch to elect an owner, broadcasts diff --git a/simpler_setup/tools/swimlane_converter.py b/simpler_setup/tools/swimlane_converter.py index eb28b5b5d5..98c1a6dd63 100644 --- a/simpler_setup/tools/swimlane_converter.py +++ b/simpler_setup/tools/swimlane_converter.py @@ -1759,6 +1759,8 @@ def sched_lane_tid(thread_idx, lane=0): "drain_prepare": "cq_build_attempt_runnable", # inner: cluster scan + build_payload "drain_publish": "cq_build_attempt_passed", # inner: MMIO write_reg per subtask (the cohort launch) "graph_prepare": "rail_animation", # bounded Scheduler-side Definition expansion + # 调用阶段:设备侧 Scheduler 已结束,Host 正在离线转换 terminal phase。 + "terminal_close": "olive", # successful-run bulk lifecycle closure # Inner in TMR; standalone on HBG's dedicated P thread. "resolve": "vsync_highlight_color", # on_task_complete: walk consumer list # Separate-lane (Worker View AICPU_N) — fallback color if it ever lands on Sched @@ -1935,6 +1937,8 @@ def _find_containing_complete(thread_idx: int, finish_us: float): "drain_prepare", "drain_publish", "graph_prepare", + # 调用阶段:设备侧 Scheduler 已结束,Host 筛选并导出已回传的 terminal phase。 + "terminal_close", ): continue start_us = record["start_time_us"] diff --git a/src/a5/runtime/tensormap_and_ringbuffer/aicpu/aicpu_executor.cpp b/src/a5/runtime/tensormap_and_ringbuffer/aicpu/aicpu_executor.cpp index 0698651a86..2139b173e1 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/aicpu/aicpu_executor.cpp +++ b/src/a5/runtime/tensormap_and_ringbuffer/aicpu/aicpu_executor.cpp @@ -304,13 +304,11 @@ int32_t AicpuExecutor::init(Runtime *runtime) { if (init_failed_.load(std::memory_order_acquire)) return -1; } #else + // 调用阶段:Orchestrator 和 Scheduler 启动前的 AICPU 握手阶段。 // Perf path: a scheduler that sees an invalid core report (its own or a - // peer's, observed so far) latches completed_ via abort_and_shutdown, which - // stops any peer still entering dispatch (run()'s is_completed() gate). A - // peer that already passed that gate is not joined here — its own cores are - // valid (it handshaked them), and the failure ends in the host device reset - // that reaps every core, so the residual overlap is bounded and - // non-corrupting. finished_count_ is reset per-run in deinit(), not here. + // peer's, observed so far) latches completed_ via abort_and_shutdown. + // Peers that have not entered dispatch stop at run()'s completion gate; + // the failure ends in the host device reset that reaps every core. if (sched_ctx_.handshake_failed()) { sched_ctx_.abort_and_shutdown(runtime); init_failed_.store(true, std::memory_order_release); @@ -851,6 +849,8 @@ int32_t AicpuExecutor::run(Runtime *runtime) { LOG_INFO("Thread %d: Orchestrator completed", thread_idx); } + // 调用阶段:Scheduler 线程进入调度;正常并发模式下 Orchestrator 此时仍可能运行。 + // 成功路径的 Scheduler 参加统一终态协议;已完成的失败路径直接进入 shutdown。 // Scheduler thread (orchestrator thread skips dispatch and exits after orchestration) if (!sched_ctx_.is_completed() && thread_idx < sched_thread_num_) { // Device orchestration: wait for the primary orchestrator to initialize the SM header diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/async_wait.h b/src/a5/runtime/tensormap_and_ringbuffer/runtime/async_wait.h index 6c282cb065..e1b88ed301 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/async_wait.h +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/async_wait.h @@ -188,6 +188,8 @@ struct AsyncWaitList { ChipTaskSlotState **deferred_release_slot_states{nullptr}; int32_t *deferred_release_count{nullptr}; int32_t deferred_release_capacity{0}; + // 调用阶段:Scheduler 运行期间;Orchestrator 可能仍在运行,也可能已经结束。 + const std::atomic *release_seal{nullptr}; int32_t inline_completed{0}; #if SIMPLER_SCHED_PROFILING int32_t thread_idx{0}; @@ -303,10 +305,11 @@ struct AsyncWaitList { } template + // 调用阶段:Scheduler 主循环轮询异步完成;Orchestrator 可能仍在运行,也可能已经结束。 AsyncPollResult poll_and_complete( AICoreCompletionMailbox *aicore_mailbox, SchedulerState *sched, ChipTaskSlotState **deferred_release_slot_states, int32_t &deferred_release_count, - int32_t deferred_release_capacity + int32_t deferred_release_capacity, const std::atomic *release_seal #if SIMPLER_SCHED_PROFILING , int thread_idx diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler.h b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler.h index ced25fe2bf..04c319aa33 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler.h +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler.h @@ -637,6 +637,45 @@ struct SchedulerState { return advanced; } + // 调用阶段:Orchestrator 已结束、所有任务已完成,并且全部 Scheduler 调度循环均已退出。 + // 仅 terminal barrier 的最后到达线程调用,统一关闭仍存活的 ring slot。 + int32_t terminal_close_live_slots() { + int32_t current_task_indices[CHIP_MAX_RING_DEPTH]; + int32_t total_closed = 0; + + for (int32_t ring_id = 0; ring_id < CHIP_MAX_RING_DEPTH; ring_id++) { + auto &ring_sched = ring_sched_states[ring_id]; + int32_t current_task_index = ring_sched.ring->fc.current_task_index.load(std::memory_order_acquire); + int32_t live_count = current_task_index - ring_sched.last_task_alive; + if (live_count < 0 || static_cast(live_count) > ring_sched.ring->task_window_size) { + LOG_ERROR( + "terminal lifecycle close has invalid ring %d interval [%d, %d) for window %" PRIu64, ring_id, + ring_sched.last_task_alive, current_task_index, ring_sched.ring->task_window_size + ); + return -1; + } + current_task_indices[ring_id] = current_task_index; + total_closed += live_count; + } + + for (int32_t ring_id = 0; ring_id < CHIP_MAX_RING_DEPTH; ring_id++) { + auto &ring_sched = ring_sched_states[ring_id]; + int32_t current_task_index = current_task_indices[ring_id]; + for (int32_t id = ring_sched.last_task_alive; id < current_task_index; id++) { + ChipTaskSlotState &slot_state = ring_sched.ring->get_slot_state_by_task_id(id); + slot_state.task_state.store(CHIP_TASK_CONSUMED, std::memory_order_relaxed); + slot_state.reset_for_reuse(); + } + ring_sched.last_task_alive = current_task_index; + ring_sched.sync_to_sm(true); + } + + advance_pending_mask.store(0, std::memory_order_relaxed); + publication_request_mask.store(0, std::memory_order_relaxed); + publication_ack_mask.store(0, std::memory_order_relaxed); + return total_closed; + } + bool try_claim_ready_once(ChipTaskSlotState &slot_state) { uint8_t flags = slot_state.lifecycle_flags.load(std::memory_order_acquire); for (;;) { @@ -1323,26 +1362,37 @@ AsyncWaitList::try_inline_complete_locked(AsyncWaitList::DrainCompletionSink &si #else sink.sched->on_task_complete(slot_state); #endif + // 调用阶段:Scheduler drain 正在处理异步完成;Orchestrator 可能仍在运行,也可能已经结束。 + // 仅当本地 release buffer 满时读取 seal,Orchestrator 结束后跳过逐 task release。 + bool release_elided = false; if (*sink.deferred_release_count >= sink.deferred_release_capacity) { - while (*sink.deferred_release_count > 0) { + release_elided = sink.release_seal != nullptr && sink.release_seal->load(std::memory_order_acquire); + if (release_elided) { + *sink.deferred_release_count = 0; + } else { + while (*sink.deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING - (void)sink.sched->on_task_release( - *sink.deferred_release_slot_states[--(*sink.deferred_release_count)], sink.thread_idx - ); + (void)sink.sched->on_task_release( + *sink.deferred_release_slot_states[--(*sink.deferred_release_count)], sink.thread_idx + ); #else - sink.sched->on_task_release(*sink.deferred_release_slot_states[--(*sink.deferred_release_count)]); + sink.sched->on_task_release(*sink.deferred_release_slot_states[--(*sink.deferred_release_count)]); #endif + } } } - sink.deferred_release_slot_states[(*sink.deferred_release_count)++] = &slot_state; + if (!release_elided) { + sink.deferred_release_slot_states[(*sink.deferred_release_count)++] = &slot_state; + } sink.inline_completed++; return true; } template +// 调用阶段:Scheduler 主循环轮询异步任务;Orchestrator 可能仍在运行,也可能已经结束。 inline AsyncPollResult AsyncWaitList::poll_and_complete( AICoreCompletionMailbox *aicore_mailbox, SchedulerState *sched, ChipTaskSlotState **deferred_release_slot_states, - int32_t &deferred_release_count, int32_t deferred_release_capacity + int32_t &deferred_release_count, int32_t deferred_release_capacity, const std::atomic *release_seal #if SIMPLER_SCHED_PROFILING , int thread_idx @@ -1356,6 +1406,8 @@ inline AsyncPollResult AsyncWaitList::poll_and_complete( sink.deferred_release_slot_states = deferred_release_slot_states; sink.deferred_release_count = &deferred_release_count; sink.deferred_release_capacity = deferred_release_capacity; + // 调用阶段:Scheduler 运行期间向 mailbox drain 传递 Orchestrator 完成标志。 + sink.release_seal = release_seal; #if SIMPLER_SCHED_PROFILING sink.thread_idx = thread_idx; #endif @@ -1394,16 +1446,28 @@ inline AsyncPollResult AsyncWaitList::poll_and_complete( #else sched->on_task_complete(*entry.slot_state); #endif + // 调用阶段:Scheduler 完成异步 task;Orchestrator 可能仍在运行,也可能已经结束。 + // seal 只在 release buffer 满的边界生效。 + bool release_elided = false; if (deferred_release_count >= deferred_release_capacity) { - while (deferred_release_count > 0) { + release_elided = release_seal != nullptr && release_seal->load(std::memory_order_acquire); + if (release_elided) { + deferred_release_count = 0; + } else { + while (deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING - (void)sched->on_task_release(*deferred_release_slot_states[--deferred_release_count], thread_idx); + (void)sched->on_task_release( + *deferred_release_slot_states[--deferred_release_count], thread_idx + ); #else - sched->on_task_release(*deferred_release_slot_states[--deferred_release_count]); + sched->on_task_release(*deferred_release_slot_states[--deferred_release_count]); #endif + } } } - deferred_release_slot_states[deferred_release_count++] = entry.slot_state; + if (!release_elided) { + deferred_release_slot_states[deferred_release_count++] = entry.slot_state; + } result.completed++; int32_t last = count - 1; diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_cold_path.cpp b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_cold_path.cpp index 8b2b613b0a..b8ecea03c4 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_cold_path.cpp +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_cold_path.cpp @@ -10,8 +10,11 @@ */ #include "scheduler_context.h" +// 使用阶段:Scheduler 调度结束后的 terminal closure 冷路径。 +#include #include #include +#include #include "common/unified_log.h" #include "aicpu/aicpu_device_config.h" @@ -52,6 +55,48 @@ static bool latch_scheduler_error(SharedMemoryHeader *header, int32_t thread_idx return won; } +// 调用阶段:当前 Scheduler 线程已经退出调度循环;Orchestrator 已结束且全部任务已完成时进入 barrier。 +// 最后到达的 Scheduler 执行 bulk closure,其他 Scheduler 等待 closure 发布完成。 +int32_t SchedulerContext::finish_successful_terminal(SharedMemoryHeader *header, int32_t thread_idx) { + bool successful_terminal = orchestrator_done_.load(std::memory_order_acquire) && total_tasks_ > 0 && + completed_tasks_.load(std::memory_order_acquire) >= total_tasks_ && + header->orch_error_code.load(std::memory_order_acquire) == SIMPLER_ERROR_NONE && + header->sched_error_code.load(std::memory_order_acquire) == SIMPLER_ERROR_NONE; + if (!successful_terminal) return 0; + + int32_t terminal_close_result = 0; + int32_t arrived = terminal_close_arrived_.fetch_add(1, std::memory_order_acq_rel) + 1; + if (arrived == active_sched_threads_) { +#if SIMPLER_DFX + uint64_t terminal_close_t0 = get_sys_cnt_aicpu(); +#endif + terminal_close_result = sched_->terminal_close_live_slots(); +#if SIMPLER_DFX + if (chip_swimlane_level_ >= ChipSwimlaneLevel::SCHED_PHASES) { + uint64_t terminal_close_t1 = get_sys_cnt_aicpu(); + int16_t phase_depth[CHIP_SWIMLANE_NUM_QUEUE_SHAPES]; + constexpr size_t kMax = static_cast(std::numeric_limits::max()); + for (int s = 0; s < CHIP_SWIMLANE_NUM_QUEUE_SHAPES; s++) { + size_t qsize = sched_->ready_queues[s].size() + sched_->ready_sync_queues[s].size(); + phase_depth[s] = static_cast(std::min(qsize, kMax)); + } + chip_swimlane_aicpu_record_sched_phase( + thread_idx, ChipSwimlaneSchedPhaseKind::TerminalClose, terminal_close_t0, terminal_close_t1, + sched_chip_swimlane_[thread_idx].sched_loop_count, + static_cast(terminal_close_result < 0 ? 0 : terminal_close_result), /*pop_hit=*/0, + /*pop_miss=*/0, phase_depth, phase_depth + ); + } +#endif + terminal_close_status_.store(terminal_close_result < 0 ? -1 : 1, std::memory_order_release); + } else { + while ((terminal_close_result = terminal_close_status_.load(std::memory_order_acquire)) == 0) { + SPIN_WAIT_HINT(); + } + } + return terminal_close_result < 0 ? -1 : 0; +} + LoopAction SchedulerContext::handle_orchestrator_exit( int32_t thread_idx, SharedMemoryHeader *header, Runtime *runtime, int32_t &task_count ) { @@ -1152,6 +1197,9 @@ int32_t SchedulerContext::pre_handshake_init( // released to dispatch. completed_tasks_.store(0, std::memory_order_release); orchestrator_done_.store(false, std::memory_order_release); + // 调用阶段:Orchestrator 与 Scheduler 启动前,由握手 leader 初始化本轮 terminal barrier。 + terminal_close_arrived_.store(0, std::memory_order_release); + terminal_close_status_.store(0, std::memory_order_release); func_id_to_addr_ = runtime->dev.func_id_to_addr_; // total_tasks_ must be read before hs_setup_done_ is published: on the @@ -1315,6 +1363,9 @@ void SchedulerContext::deinit() { total_tasks_ = 0; orchestrator_done_.store(false, std::memory_order_release); completed_.store(false, std::memory_order_release); + // 调用阶段:Orchestrator 与全部 Scheduler 均已结束,deinit 为下一轮清理终态协调状态。 + terminal_close_arrived_.store(0, std::memory_order_release); + terminal_close_status_.store(0, std::memory_order_release); // Reset core discovery and assignment state aic_count_ = 0; diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_completion.cpp b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_completion.cpp index f7a054f818..4ce6f01a29 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_completion.cpp +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_completion.cpp @@ -198,8 +198,12 @@ void SchedulerContext::complete_slot_task( } chip_swimlane.phase_complete_count++; #endif + // 调用阶段:Scheduler 处理 AICore completion;Orchestrator 可能仍在运行,也可能已经结束。 + // 仅在本地 release buffer 满时读取 seal;编排结束后丢弃 release 债务,稍后统一收口。 if (deferred_release_count < DEFERRED_RELEASE_CAP) { deferred_release_slot_states[deferred_release_count++] = &slot_state; + } else if (orchestrator_done_.load(std::memory_order_acquire)) { + deferred_release_count = 0; } else { while (deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_context.h b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_context.h index a6c7106edd..88fdf7ccfc 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_context.h +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_context.h @@ -214,6 +214,12 @@ class SchedulerContext { // Platform AICore-register base array (set by AicpuExecutor before init()). uint64_t regs_{0}; + // 使用阶段:Orchestrator 已结束、Scheduler 调度循环退出后的成功终态协调。 + // status: 0 表示 closure 未完成,1 表示成功,-1 表示失败。两个 32 位原子 + // 复用 regs_ 后的 8 字节尾部空洞,不扩大 SchedulerContext。 + std::atomic terminal_close_arrived_{0}; + std::atomic terminal_close_status_{0}; + // ========================================================================= // Core management (scheduler_cold_path.cpp) // ========================================================================= @@ -470,6 +476,9 @@ class SchedulerContext { #endif ); + // 调用阶段:单个 Scheduler 调度循环退出后;最后到达线程负责 terminal closure。 + __attribute__((noinline, cold)) int32_t finish_successful_terminal(SharedMemoryHeader *header, int32_t thread_idx); + #if SIMPLER_DFX __attribute__((noinline, cold)) void log_chip_swimlane_summary(int32_t thread_idx, int32_t cur_thread_completed); #endif diff --git a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_dispatch.cpp b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_dispatch.cpp index 6f2cf594c1..eeddf6080e 100644 --- a/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_dispatch.cpp +++ b/src/a5/runtime/tensormap_and_ringbuffer/runtime/scheduler/scheduler_dispatch.cpp @@ -1044,8 +1044,10 @@ int32_t SchedulerContext::resolve_and_dispatch(Runtime *runtime, int32_t thread_ if (rt_ != nullptr && rt_->aicore_mailbox != nullptr && (sched_->async_wait_list.count > 0 || rt_->aicore_mailbox->has_pending())) { + // 调用阶段:Scheduler 主循环轮询异步完成;Orchestrator 可能仍在运行,也可能已经结束。 AsyncPollResult poll_result = sched_->async_wait_list.poll_and_complete( - rt_->aicore_mailbox, sched_, deferred_release_slot_states, deferred_release_count, DEFERRED_RELEASE_CAP + rt_->aicore_mailbox, sched_, deferred_release_slot_states, deferred_release_count, DEFERRED_RELEASE_CAP, + &orchestrator_done_ #if SIMPLER_SCHED_PROFILING , thread_idx @@ -1174,18 +1176,24 @@ int32_t SchedulerContext::resolve_and_dispatch(Runtime *runtime, int32_t thread_ } #endif // Dummy tasks have no subtasks to retire and no fanout pre-conditions - // beyond their own producers; release self-reference so the slot can - // reach CONSUMED once all consumers drain. + // beyond their own producers. While lifecycle reclamation is active, + // release their self-reference so the slot can reach CONSUMED. + // 调用阶段:Scheduler 主循环完成 dummy task;Orchestrator 可能仍在运行,也可能已经结束。 + // seal 只在本地 release buffer 满时生效。 deferred_release_slot_states[deferred_release_count++] = &dummy_slot; if (deferred_release_count >= DEFERRED_RELEASE_CAP) { - while (deferred_release_count > 0) { + if (orchestrator_done_.load(std::memory_order_acquire)) { + deferred_release_count = 0; + } else { + while (deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING - (void)sched_->on_task_release( - *deferred_release_slot_states[--deferred_release_count], thread_idx - ); + (void)sched_->on_task_release( + *deferred_release_slot_states[--deferred_release_count], thread_idx + ); #else - sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count]); + sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count]); #endif + } } } int32_t prev = completed_tasks_.fetch_add(1, std::memory_order_relaxed); @@ -1298,15 +1306,21 @@ int32_t SchedulerContext::resolve_and_dispatch(Runtime *runtime, int32_t thread_ 0; uint32_t released_count = static_cast(deferred_release_count); #endif - while (deferred_release_count > 0) { + // 调用阶段:Scheduler 主循环本轮无 completion/dispatch 进展,准备排空本线程 backlog。 + // Orchestrator 运行时正常 release;Orchestrator 结束后跳过,等待 terminal closure。 + bool release_elided = deferred_release_count > 0 && orchestrator_done_.load(std::memory_order_acquire); + while (!release_elided && deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING (void)sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count], thread_idx); #else sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count]); #endif } + if (release_elided) { + deferred_release_count = 0; + } #if SIMPLER_DFX - if (release_t0 != 0) { + if (release_t0 != 0 && !release_elided) { chip_swimlane_aicpu_record_sched_phase( thread_idx, ChipSwimlaneSchedPhaseKind::Release, release_t0, get_sys_cnt_aicpu(), chip_swimlane.sched_loop_count, released_count @@ -1384,18 +1398,22 @@ int32_t SchedulerContext::resolve_and_dispatch(Runtime *runtime, int32_t thread_ } } - // Drain any entries left in the deferred-release batch. The in-loop flush - // only fires on idle iterations and on buffer-full; a loop exit while the - // last iteration made progress can leave entries un-released. Drop them - // here so every consumed producer slot completes its on_task_release - // regardless of which loop-exit path fired. - while (deferred_release_count > 0) { + // 调用阶段:当前 Scheduler 调度循环已经结束,其他 Scheduler 可能仍在退出途中。 + // Orchestrator 未结束时继续精确 release;Orchestrator 已结束时把剩余 slot 交给 terminal closure。 + // An unsealed graph still needs exact lifecycle bookkeeping on every exit. + // Once orchestration is sealed, terminal closure owns the remaining slots. + bool release_elided = deferred_release_count > 0 && orchestrator_done_.load(std::memory_order_acquire); + while (!release_elided && deferred_release_count > 0) { #if SIMPLER_SCHED_PROFILING (void)sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count], thread_idx); #else sched_->on_task_release(*deferred_release_slot_states[--deferred_release_count]); #endif } + // 调用阶段:当前 Scheduler 已结束调度;在此等待全部 Scheduler 到齐并执行成功终态收口。 + if (finish_successful_terminal(header, thread_idx) < 0) { + timeout_rc = -1; + } #if SIMPLER_DFX // Final-drain: emit any pop_hit / pop_miss accrued since the last diff --git a/src/common/platform/include/common/chip_swimlane_profiling.h b/src/common/platform/include/common/chip_swimlane_profiling.h index bb57f6f560..bb40663e1c 100644 --- a/src/common/platform/include/common/chip_swimlane_profiling.h +++ b/src/common/platform/include/common/chip_swimlane_profiling.h @@ -550,6 +550,11 @@ enum class ChipSwimlaneSchedPhaseKind : uint32_t { // phase_data.graph_task identifies the ring-0 outer Graph task and // tasks_processed is the number of in-graph tasks patched in this slice. GraphPrepare = 13, + // Outer (sched lane): successful-run bulk lifecycle closure after every + // scheduler thread has left dispatch/completion. tasks_processed is the + // number of live slots terminalized across all rings. + // 记录阶段:Orchestrator 已结束且全部 Scheduler 调度循环退出后,由 terminal leader 写入。 + TerminalClose = 14, }; /** Index layout of the queue-depth snapshot arrays below: AIC=0, AIV=1, MIX=2. diff --git a/src/common/platform/shared/host/chip_swimlane_collector.cpp b/src/common/platform/shared/host/chip_swimlane_collector.cpp index 58fdcbb2f1..ce0061c1fc 100644 --- a/src/common/platform/shared/host/chip_swimlane_collector.cpp +++ b/src/common/platform/shared/host/chip_swimlane_collector.cpp @@ -1171,6 +1171,9 @@ int ChipSwimlaneCollector::export_swimlane_json() { return "async_poll"; case ChipSwimlaneSchedPhaseKind::GraphPrepare: return "graph_prepare"; + // 调用阶段:设备侧 Orchestrator/Scheduler 均已结束,Host 导出已采集的 terminal phase。 + case ChipSwimlaneSchedPhaseKind::TerminalClose: + return "terminal_close"; } return "unknown"; }; diff --git a/tests/ut/cpp/a5/test_scheduler_state.cpp b/tests/ut/cpp/a5/test_scheduler_state.cpp index 2cb54304d0..58dea0a449 100644 --- a/tests/ut/cpp/a5/test_scheduler_state.cpp +++ b/tests/ut/cpp/a5/test_scheduler_state.cpp @@ -195,6 +195,99 @@ TEST_F(SchedulerStateTest, ConsumedTransition) { EXPECT_EQ(slot.task_state.load(), CHIP_TASK_CONSUMED); } +// 测试阶段:离线单元测试,模拟 Orchestrator 与全部 Scheduler 结束后的 terminal closure。 +TEST_F(SchedulerStateTest, TerminalCloseResetsLiveIntervalAndPublishesTail) { + constexpr int32_t ring_id = 2; + SharedMemoryRingHeader &ring = sm_handle->header->rings[ring_id]; + SchedulerState::RingSchedState &ring_sched = sched.ring_sched_states[ring_id]; + + ring.fc.current_task_index.store(5, std::memory_order_release); + ring.fc.last_task_alive.store(1, std::memory_order_release); + ring_sched.last_task_alive = 1; + ring_sched.last_published_to_sm = 1; + for (int32_t task_id = 1; task_id < 5; task_id++) { + init_ring_slot(ring, task_id, CHIP_TASK_COMPLETED, ring_id, 7, 2); + ChipTaskSlotState &slot = ring.get_slot_state_by_task_id(task_id); + slot.fanin_refcount.store(3, std::memory_order_relaxed); + slot.completed_subtasks.store(1, std::memory_order_relaxed); + } + + EXPECT_EQ(sched.terminal_close_live_slots(), 4); + EXPECT_EQ(ring_sched.last_task_alive, 5); + EXPECT_EQ(ring_sched.last_published_to_sm, 5); + EXPECT_EQ(ring.fc.last_task_alive.load(std::memory_order_acquire), 5); + for (int32_t task_id = 1; task_id < 5; task_id++) { + ChipTaskSlotState &slot = ring.get_slot_state_by_task_id(task_id); + EXPECT_EQ(slot.task_state.load(std::memory_order_relaxed), CHIP_TASK_CONSUMED); + EXPECT_EQ(slot.fanin_refcount.load(std::memory_order_relaxed), 0); + EXPECT_EQ(slot.fanout_refcount.load(std::memory_order_relaxed), 0u); + EXPECT_EQ(slot.fanout_count, FANOUT_SCOPE_BIT); + EXPECT_EQ(slot.completed_subtasks.load(std::memory_order_relaxed), 0); + } +} + +// 测试阶段:离线单元测试,覆盖 Scheduler 结束后 terminal closure 的非法 ring 区间。 +TEST_F(SchedulerStateTest, TerminalCloseRejectsIntervalLargerThanWindow) { + constexpr int32_t ring_id = 1; + SharedMemoryRingHeader &ring = sm_handle->header->rings[ring_id]; + SchedulerState::RingSchedState &ring_sched = sched.ring_sched_states[ring_id]; + + int32_t invalid_head = static_cast(ring.task_window_size) + 1; + ring.fc.current_task_index.store(invalid_head, std::memory_order_release); + ring.fc.last_task_alive.store(0, std::memory_order_release); + ring_sched.last_task_alive = 0; + + EXPECT_EQ(sched.terminal_close_live_slots(), -1); + EXPECT_EQ(ring_sched.last_task_alive, 0); + EXPECT_EQ(ring.fc.last_task_alive.load(std::memory_order_acquire), 0); +} + +// 测试阶段:离线单元测试,模拟 Orchestrator 运行期间的 async completion。 +TEST_F(SchedulerStateTest, AsyncInlineCompletionDefersReleaseBeforeOrchestratorDone) { + alignas(64) ChipTaskSlotState slot; + init_slot(slot, CHIP_TASK_PENDING, 0, 1); + ChipTaskSlotState *deferred[1]{}; + int32_t deferred_count = 0; + std::atomic release_seal{false}; + AsyncWaitList::DrainCompletionSink sink{}; + sink.sched = &sched; + sink.deferred_release_slot_states = deferred; + sink.deferred_release_count = &deferred_count; + sink.deferred_release_capacity = 1; + sink.release_seal = &release_seal; + + EXPECT_TRUE(sched.async_wait_list.try_inline_complete_locked(sink, slot)); + EXPECT_EQ(slot.task_state.load(), CHIP_TASK_COMPLETED); + EXPECT_EQ(deferred_count, 1); + EXPECT_EQ(deferred[0], &slot); +} + +// 测试阶段:离线单元测试,模拟 Orchestrator 结束后、Scheduler 尚未结束时的 async completion。 +TEST_F(SchedulerStateTest, AsyncInlineCompletionElidesReleaseAtSealedCapacityBoundary) { + alignas(64) ChipTaskSlotState first; + alignas(64) ChipTaskSlotState second; + init_slot(first, CHIP_TASK_PENDING, 0, 1); + init_slot(second, CHIP_TASK_PENDING, 0, 1); + ChipTaskSlotState *deferred[1]{}; + int32_t deferred_count = 0; + std::atomic release_seal{true}; + AsyncWaitList::DrainCompletionSink sink{}; + sink.sched = &sched; + sink.deferred_release_slot_states = deferred; + sink.deferred_release_count = &deferred_count; + sink.deferred_release_capacity = 1; + sink.release_seal = &release_seal; + + EXPECT_TRUE(sched.async_wait_list.try_inline_complete_locked(sink, first)); + EXPECT_EQ(first.task_state.load(), CHIP_TASK_COMPLETED); + EXPECT_EQ(deferred_count, 1); + EXPECT_EQ(deferred[0], &first); + + EXPECT_TRUE(sched.async_wait_list.try_inline_complete_locked(sink, second)); + EXPECT_EQ(second.task_state.load(), CHIP_TASK_COMPLETED); + EXPECT_EQ(deferred_count, 0); +} + TEST_F(SchedulerStateTest, ConsumedHeadAdvancesAfterContendedAdvanceLock) { constexpr int32_t ring_id = CHIP_MAX_RING_DEPTH - 1; constexpr int32_t head_task_id = 0;