diff --git a/vnt-core/src/utils/task_control.rs b/vnt-core/src/utils/task_control.rs index bed1d28..d58a325 100644 --- a/vnt-core/src/utils/task_control.rs +++ b/vnt-core/src/utils/task_control.rs @@ -155,10 +155,14 @@ impl TaskGroup { pub async fn wait_all_stopped(&self) { loop { + // 先注册等待再检查条件,避免在检查与等待之间丢失唤醒 + let notified = self.inner.all_stopped_notify.notified(); + tokio::pin!(notified); + notified.as_mut().enable(); if self.inner.all_tasks_stopped() { return; } - self.inner.all_stopped_notify.notified().await; + notified.await; } } } @@ -253,3 +257,33 @@ impl Drop for TaskGroupGuard { } } } + +#[cfg(test)] +mod tests { + use super::*; + + /// 所有任务自然结束后 wait_all_stopped 必须返回。 + /// 覆盖两个关键点:任务自然耗尽时 remove_task 置 stopped 并唤醒; + /// 等待方先注册再检查,不会因竞态错过唤醒而永久挂起。 + #[tokio::test] + async fn test_wait_all_stopped_after_natural_completion() { + let manager = TaskGroupManager::new(); + let (group, _guard) = manager.create_task().unwrap(); + + let waiter = { + let group = group.clone(); + tokio::spawn(async move { group.wait_all_stopped().await }) + }; + // 让 waiter 先进入等待 + tokio::task::yield_now().await; + + let _sub = group.spawn(async { + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + }); + + tokio::time::timeout(std::time::Duration::from_secs(2), waiter) + .await + .expect("wait_all_stopped should return after all tasks complete") + .unwrap(); + } +} diff --git a/vnt-web/src/service_http.rs b/vnt-web/src/service_http.rs index 2c23a08..e898815 100644 --- a/vnt-web/src/service_http.rs +++ b/vnt-web/src/service_http.rs @@ -36,7 +36,7 @@ use vnt_core::utils::task_control::TaskGroupManager; const CONFIG_DIR: &str = "vnt_config"; const CURRENT_CONFIG_RECORD: &str = "vnt_current_config.txt"; -#[derive(Serialize, Clone, Copy, PartialEq, Eq, Default)] +#[derive(Serialize, Clone, Copy, PartialEq, Eq, Default, Debug)] #[serde(rename_all = "lowercase")] enum VntStatus { #[default] @@ -56,6 +56,8 @@ struct HttpAppStateInner { vnt: Option, status: VntStatus, start_logs: Vec, + /// 启动任务句柄,用于在 Starting 状态中断注册重试循环 + start_handle: Option>, } impl HttpAppState { @@ -120,6 +122,13 @@ impl HttpAppState { self.inner.lock().status } + /// 中断启动任务(如注册重试循环)。任务已完成时为空操作。 + fn abort_start_task(&self) { + if let Some(handle) = self.inner.lock().start_handle.take() { + handle.abort(); + } + } + fn timestamp() -> String { let now = OffsetDateTime::now_local().unwrap_or_else(|_| OffsetDateTime::now_utc()); let format = format_description::parse("[hour]:[minute]:[second]").unwrap(); @@ -522,7 +531,7 @@ async fn start_vnt_internal( state.record_log("创建组网管理器"); let state_clone = state.clone(); - tokio::spawn(async move { + let start_handle = tokio::spawn(async move { let result = start_vnt_network( state_clone.clone(), file_name, @@ -540,6 +549,7 @@ async fn start_vnt_internal( } drop(on_error_guard); }); + state.inner.lock().start_handle = Some(start_handle); Ok(()) } @@ -636,8 +646,10 @@ async fn start_vnt_network( state.starting_to_running(); - // 启动网络管理任务 - task_group.spawn(async move { + // 启动网络管理任务。 + // 注意必须在任务组外等待:等待目标就是这个 task_group, + // 若 spawn 进组内会形成自引用等待,网络自行停止时永不返回 + tokio::spawn(async move { network_manager.wait_all_stopped().await; drop(task_group_guard); drop(network_manager); @@ -682,6 +694,8 @@ async fn stop_vnt_handler(State(state): State) -> Json