Rust状态机异步触发回调遇E0521编译错误求助
Rust状态机异步触发回调的生命周期问题
问题背景
我在Rust中实现一个状态机,状态使用泛型T适配多种场景,用Vec存储所有注册的回调。最初未使用生命周期参数时出现生命周期问题,参考资料添加生命周期'a后,同步触发回调可正常运行,但在spawn线程中触发回调时遇到编译错误。
状态机结构体定义
pub struct StateMachine<'a, T> where T:Clone+Eq+'a { state: RwLock<T>, listeners2: Vec<Arc<Mutex<ListenerCallback<'a, T>>>>, } pub type ListenerCallback<'a, T> = dyn FnMut(T) -> Result<()> + Send + Sync + 'a ;
状态变更触发回调实现
pub async fn try_set(&mut self, new_state:T) -> Result<()> { if (block_on(self.state.read()).deref().eq(&new_state)) { return Ok(()) } // todo change the state // fire every listener in spawn let mut fire_results = vec![]; for listener in &mut self.listeners2 { let state = new_state.clone(); let fire_listener = listener.clone(); fire_results.push(tokio::spawn(async move { let mut guard = fire_listener.lock().unwrap(); guard.deref_mut()(state); })); } // if fire result return Err, return it for fire_result in fire_results { fire_result.await?; } Ok(()) }
编译错误信息
error[E0521]: borrowed data escapes outside of associated function --> src/taf/taf-core/src/execution/state_machine.rs:54:33 | 15 | impl<'a,T> StateMachine<'a,T> where T:Clone+Eq+Send { | -- lifetime `'a` defined here ... 34 | pub async fn try_set(&mut self, new_state:T) -> Result<()> { | --------- `self` is a reference that is only valid in the associated function body ... 54 | let fire_listener = listener.clone(); | ^^^^^^^^^^^^^^^^ | | | `self` escapes the associated function body here | argument requires that `'a` must outlive `'static`
复现Demo
可正常运行的同步Demo
use std::ops::DerefMut; use std::sync::{Arc, Mutex, RwLock}; use anyhow::Result; use dashmap::DashMap; struct StateMachine<'a,T> where T:Clone+Eq+'a { state: T, listeners: Vec<Box<Callback<'a, T>>>, } type Callback<'a, T> = dyn FnMut(T) -> Result<()> + Send + Sync + 'a; impl<'a, T> StateMachine<'a,T> where T:Clone+Eq+'a { pub fn new(init_state: T) -> Self { StateMachine { state: init_state, listeners: vec![] } } pub fn add_listener(&mut self, listener: Box<Callback<'a, T>>) -> Result<()> { self.listeners.push(listener); Ok(()) } pub fn set(&mut self, new_state: T) -> Result<()> { self.state = new_state.clone(); for listener in &mut self.listeners { listener(new_state.clone()); } Ok(()) } } #[derive(Clone, Eq, PartialEq, Hash)] enum ExeState { Waiting, Running, Finished, Failed, } struct Execution<'a> { exec_id: String, pub state_machine: StateMachine<'a, ExeState>, } struct ExecManager<'a> { all_jobs: Arc<RwLock<DashMap<String, Execution<'a>>>>, finished_jobs: Arc<RwLock<Vec<String>>>, } impl<'a> ExecManager<'a> { pub fn new() -> Self { ExecManager { all_jobs: Arc::new(RwLock::new(DashMap::new())), finished_jobs: Arc::new(RwLock::new(vec![])) } } fn add_job(&mut self, job_id: String) { let mut execution = Execution { exec_id: job_id.clone(), state_machine: StateMachine::new(ExeState::Waiting) }; // add listener let callback_finished_jobs = self.finished_jobs.clone(); let callback_job_id = job_id.clone(); execution.state_machine.add_listener( Box::new(move |new_state| { println!("listener fired!, job_id {}", callback_job_id.clone()); if new_state == ExeState::Finished || new_state == ExeState::Failed { let mut guard = callback_finished_jobs.write().unwrap(); guard.deref_mut().push(callback_job_id.clone()); } Ok(()) })); let mut guard = self.all_jobs.write().unwrap(); guard.deref_mut().insert(job_id, execution); } fn mock_exec(&mut self, job_id: String) { let mut guard = self.all_jobs.write().unwrap(); let mut exec = guard.deref_mut().get_mut(&job_id).unwrap(); exec.state_machine.set(ExeState::Finished); } } #[test] fn test() { let mut manager = ExecManager::new(); manager.add_job(String::from("job_id1")); manager.add_job(String::from("job_id2")); manager.mock_exec(String::from("job_id1")); manager.mock_exec(String::from("job_id2")); }
触发错误的异步Demo
use std::ops::DerefMut; use std::sync::{Arc, Mutex, RwLock}; use anyhow::Result; use dashmap::DashMap; struct StateMachine<'a,T> where T:Clone+Eq+Send+'a { state: T, listeners: Vec<Arc<Mutex<Box<Callback<'a, T>>>>>, } type Callback<'a, T> = dyn FnMut(T) -> Result<()> + Send + Sync + 'a; impl<'a, T> StateMachine<'a,T> where T:Clone+Eq+Send+'a { pub fn new(init_state: T) -> Self { StateMachine { state: init_state, listeners: vec![] } } pub fn add_listener(&mut self, listener: Box<Callback<'a, T>>) -> Result<()> { self.listeners.push(Arc::new(Mutex::new(listener))); Ok(()) } pub fn set(&mut self, new_state: T) -> Result<()> { self.state = new_state.clone(); for listener in &mut self.listeners { let spawn_listener = listener.clone(); tokio::spawn(async move { let mut guard = spawn_listener.lock().unwrap(); guard.deref_mut()(new_state.clone()); }); } Ok(()) } } #[derive(Clone, Eq, PartialEq, Hash)] enum ExeState { Waiting, Running, Finished, Failed, } struct Execution<'a> { exec_id: String, pub state_machine: StateMachine<'a, ExeState>, } struct ExecManager<'a> { all_jobs: Arc<RwLock<DashMap<String, Execution<'a>>>>, finished_jobs: Arc<RwLock<Vec<String>>>, } impl<'a> ExecManager<'a> { pub fn new() -> Self { ExecManager { all_jobs: Arc::new(RwLock::new(DashMap::new())), finished_jobs: Arc::new(RwLock::new(vec![])) } } fn add_job(&mut self, job_id: String) { let mut execution = Execution { exec_id: job_id.clone(), state_machine: StateMachine::new(ExeState::Waiting) }; // add listener let callback_finished_jobs = self.finished_jobs.clone(); let callback_job_id = job_id.clone(); execution.state_machine.add_listener( Box::new(move |new_state| { println!("listener fired!, job_id {}", callback_job_id.clone()); if new_state == ExeState::Finished || new_state == ExeState::Failed { let mut guard = callback_finished_jobs.write().unwrap(); guard.deref_mut().push(callback_job_id.clone()); } Ok(()) })); let mut guard = self.all_jobs.write().unwrap(); guard.deref_mut().insert(job_id, execution); } fn mock_exec(&mut self, job_id: String) { let mut guard = self.all_jobs.write().unwrap(); let mut exec = guard.deref_mut().get_mut(&job_id).unwrap(); exec.state_machine.set(ExeState::Finished); } } #[test] fn test() { let mut manager = ExecManager::new(); manager.add_job(String::from("job_id1")); manager.add_job(String::from("job_id2")); manager.mock_exec(String::from("job_id1")); manager.mock_exec(String::from("job_id2")); }
异步Demo编译错误
error[E0521]: borrowed data escapes outside of associated function --> generic/src/callback2.rs:34:34 | 15 | impl<'a, T> StateMachine<'a,T> where T:Clone+Eq+Send+'a { | -- lifetime `'a` defined here ... 29 | pub fn set(&mut self, new_state: T) -> Result<()> { | --------- `self` is a reference that is only valid in the associated function body ... 34 | let spawn_listener = listener.clone(); | ^^^^^^^^^^^^^^^^ | | | `self` escapes the associated function body here | argument requires that `'a` must outlive `'static` | = note: requirement occurs because of the type `std::sync::Mutex<Box<dyn FnMut(T) -> Result<(), anyhow::Error> + Send + Sync>>`, which makes the generic argument `Box<dyn FnMut(T) -> Result<(), anyhow::Error> + Send + Sync>` invariant = note: the struct `std::sync::Mutex<T>` is invariant over the parameter `T`
问题原因与解决方案
核心原因
tokio::spawn要求传入的异步任务必须满足'static生命周期——任务可能在原函数返回后仍持续运行,因此不能持有任何仅在函数体内有效的引用。而当前回调的生命周期'a绑定到了StateMachine,意味着回调可能持有StateMachine或其他外部对象的短期引用,编译器判定这些引用会随函数结束失效,因此报错self逃逸。
另外,Mutex是不变类型,不会自动协变调整生命周期,进一步放大了生命周期不匹配的问题。
解决方案
要让回调能被异步任务安全执行,核心是确保回调不依赖非'static引用,具体有两种可行方案:
方案1:移除生命周期参数,使用Arc封装回调依赖
让回调捕获的所有外部对象都通过Arc(或Arc<Mutex>)持有,使回调本身成为'static类型,无需绑定生命周期。修改后的核心代码:
pub struct StateMachine<T> where T: Clone + Eq + Send + 'static { state: RwLock<T>, listeners2: Vec<Arc<Mutex<ListenerCallback<T>>>>, } pub type ListenerCallback<T> = dyn FnMut(T) -> Result<()> + Send + Sync + 'static;
这种方式下,回调捕获的外部数据都通过Arc克隆持有,确保回调可以独立于原StateMachine存在,安全被tokio::spawn调用。
方案2:将StateMachine本身转为Arc持有
如果必须保留生命周期,可将StateMachine用Arc包裹,确保其生命周期足够覆盖异步任务的运行周期。但这种方式会增加代码复杂度,一般推荐方案1。
修改后的可运行异步Demo
use std::ops::DerefMut; use std::sync::{Arc, Mutex, RwLock}; use anyhow::Result; use dashmap::DashMap; struct StateMachine<T> where T: Clone + Eq + Send + 'static { state: T, listeners: Vec<Arc<Mutex<Box<Callback<T>>>>>, } type Callback<T> = dyn FnMut(T) -> Result<()> + Send + Sync + 'static; impl<T> StateMachine<T> where T: Clone + Eq + Send + 'static { pub fn new(init_state: T) -> Self { StateMachine { state: init_state, listeners: vec![] } } pub fn add_listener(&mut self, listener: Box<Callback<T>>) -> Result<()> { self.listeners.push(Arc::new(Mutex::new(listener))); Ok(()) } pub fn set(&mut self, new_state: T) -> Result<()> { self.state = new_state.clone(); for listener in &mut self.listeners { let spawn_listener = listener.clone(); let state_clone = new_state.clone(); tokio::spawn(async move { let mut guard = spawn_listener.lock().unwrap(); guard.deref_mut()(state_clone); }); } Ok(()) } } #[derive(Clone, Eq, PartialEq, Hash)] enum ExeState { Waiting, Running, Finished, Failed, } struct Execution { exec_id: String, pub state_machine: StateMachine<ExeState>, } struct ExecManager { all_jobs: Arc<RwLock<DashMap<String, Execution>>>, finished_jobs: Arc<RwLock<Vec<String>>>, } impl ExecManager { pub fn new() -> Self { ExecManager { all_jobs: Arc::new(RwLock::new(DashMap::new())), finished_jobs: Arc::new(RwLock::new(vec![])) } } fn add_job(&mut self, job_id: String) { let mut execution = Execution { exec_id: job_id.clone(), state_machine: StateMachine::new(ExeState::Waiting) }; // add listener let callback_finished_jobs = self.finished_jobs.clone(); let callback_job_id = job_id.clone(); execution.state_machine.add_listener( Box::new(move |new_state| { println!("listener fired!, job_id {}", callback_job_id.clone()); if new_state == ExeState::Finished || new_state == ExeState::Failed { let mut guard = callback_finished_jobs.write().unwrap(); guard.deref_mut().push(callback_job_id.clone()); } Ok(()) })); let mut guard = self.all_jobs.write().unwrap(); guard.deref_mut().insert(job_id, execution); } fn mock_exec(&mut self, job_id: String) { let mut guard = self.all_jobs.write().unwrap(); let mut exec = guard.deref_mut().get_mut(&job_id).unwrap(); exec.state_machine.set(ExeState::Finished); } } #[tokio::test] async fn test() { let mut manager = ExecManager::new(); manager.add_job(String::from("job_id1")); manager.add_job(String::from("job_id2")); manager.mock_exec(String::from("job_id1")); manager.mock_exec(String::from("job_id2")); // 等待异步任务完成 tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; }
内容的提问来源于stack exchange,提问作者lulijun
相关产品推荐
相关产品推荐

