如何在Rust中创建异步事件发射器?Discord API开发遇运行时错误
我正在尝试开发Discord API,借此学习Rust。和之前学过的语言不同,Rust的特性差异很大,现在遇到了这个错误:
thread 'main' panicked at 'Cannot start a runtime from within a runtime. This happens because a function (like
block_on) attempted to block the current thread while the thread is being used to drive asynchronous tasks.
我的代码如下:
client.rs
pub struct Client { pub listeners: Mutex<HashMap<String, Vec<Box<dyn Fn(&str) + Send + Sync>>>>, } impl Client { pub async fn on<F>(&self, event: &str, listener: F) where F: Fn(&str) + 'static + Send + Sync, { let mut listeners = self.listeners.lock().await; listeners .entry(event.to_string()) .or_insert_with(Vec::new) .push(Box::new(listener)); } pub async fn emit(&self, event: &str, data: &str) { if let Some(listeners) = self.listeners.lock().await.get(event) { for listener in listeners { listener(data); } } } pub async fn socket(&mut self) { env::set_var("RUST_BACKTRACE", "full"); /* some looping to run websocket */ match t { "READY" => {self.emit("READY", "Bot is ready")}, _ => {} } } }
main.rs
async fn main() { let mut bot: Client = Client::new(""); bot.socket(); bot.on("READY", |data| println!("event: {}", data)); }
请问如何解决该问题并正确创建异步事件发射器?
1. 修复异步调用的基础问题
当前main函数中调用bot.socket()和bot.on()时未添加.await,导致异步任务未被正确调度。更关键的是,若socket方法内部使用了block_on这类阻塞式调用,会触发"在运行时内启动新运行时"的错误——异步上下文里不能嵌套创建运行时。
2. 调整socket方法的异步逻辑
确保websocket循环完全异步,用异步框架(如tokio)的loop配合.await处理消息,禁止在异步代码中调用阻塞式同步方法。比如将websocket的接收、处理逻辑全部改为异步调用。
3. 重构main函数的异步流程
先完成监听器注册,再将长期运行的socket任务放到后台并发执行,避免阻塞主线程。示例修改如下:
use tokio; #[tokio::main] async fn main() { let mut bot: Client = Client::new(""); // 先注册事件监听器 bot.on("READY", |data| println!("event: {}", data)).await; // 将socket任务后台异步执行 tokio::spawn(async move { bot.socket().await; }); // 等待用户中断信号,防止程序直接退出 tokio::signal::ctrl_c().await.unwrap(); }
4. 优化事件发射器设计
- 若后续需要支持异步监听器,可将闭包改为异步类型,调用时添加
.await; - 替换标准库
Mutex为Tokio提供的异步互斥锁,避免异步上下文里的阻塞问题。示例修改如下:
use tokio::sync::Mutex; use std::collections::HashMap; use std::pin::Pin; use std::future::Future; pub struct Client { pub listeners: Mutex<HashMap<String, Vec<Box<dyn Fn(&str) -> Pin<Box<dyn Future<Output = ()>>> + Send + Sync + 'static>>>>, } impl Client { pub async fn on<F, Fut>(&self, event: &str, listener: F) where F: Fn(&str) -> Fut + Send + Sync + 'static, Fut: Future<Output = ()> + Send + 'static, { let mut listeners = self.listeners.lock().await; listeners .entry(event.to_string()) .or_insert_with(Vec::new) .push(Box::new(move |s| Box::pin(listener(s)))); } pub async fn emit(&self, event: &str, data: &str) { let listeners = self.listeners.lock().await; if let Some(listeners) = listeners.get(event) { for listener in listeners { listener(data).await; } } } pub async fn socket(&mut self) { env::set_var("RUST_BACKTRACE", "full"); // 异步websocket循环示例 loop { // 模拟接收READY事件 let t = "READY"; match t { "READY" => self.emit("READY", "Bot is ready").await, _ => {} } // 避免死循环占用CPU tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; } } }
5. 错误根源总结
Cannot start a runtime from within a runtime错误的核心是:在已有异步运行时的线程中,调用了会创建新运行时的阻塞方法。解决关键是全程使用异步逻辑,所有等待操作都用.await处理,不在异步上下文里调用阻塞式同步代码。
内容的提问来源于stack exchange,提问作者Gaurish Trivedi

