迷你的异步运行时
一、先搞懂几个「外部方法/类型」到底是什么
你提到的"引入了许多外部方法",其实就是这几个标准库的东西,我用大白话解释:
1. Future 特征(trait)
trait Future {
type Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}
- 一个
Future就是「一个未来会完成的任务」。 - 它只有一个核心方法
poll,返回两种结果:Poll::Ready(值)—— 任务做完了,把结果交出来。Poll::Pending—— 还没做完,下次再问。
- 关键点:
poll不是阻塞等待,而是"立刻看一眼,做完就返回 Ready,没做完就返回 Pending 马上走人"。
2. Waker 和 Context
Context就是poll时塞给你的一个"联络工具包",里面有个waker()。Waker(唤醒器)的作用是:告诉运行时"我准备好了,请再来 poll 我一次"。- 想象你在排队办事,拿了号。号没到你就坐在旁边等(返回
Pending),轮到你时广播叫号(wake),你再过去办理(再次poll)。
3. Pin<&mut Self>
- 你可以先简单理解成
&mut Self的"加强版",用来保证这个 Future 不会在内存里被乱动(移动)。初学阶段先不深究,只要知道self是固定的引用即可。
4. Arc / Mutex
Arc:多线程共享所有权(引用计数指针)。Mutex:互斥锁,保证同一时间只有一个线程能改里面的数据。Arc<Mutex<...>>组合 = "多个线程都能访问、且安全地读写同一块数据"。
5. sync_channel(同步通道)
- 一个"管道",
send往里塞东西,recv从里往外取东西。 - 本代码用它当任务队列:
Spawner往队列塞任务,Executor从队列取任务执行。
6. BoxFuture / .boxed() / waker_ref / ArcWake
BoxFuture<'static, ()>=Pin<Box<dyn Future<Output=()> + Send>>,就是把 Future 装箱、动态分发。.boxed()来自futurescrate 的扩展 trait,把任意 Future 变成BoxFuture。ArcWake是futurescrate 的 trait,让你能把Arc<Task>变成一个Waker。waker_ref(&task)基于ArcWake生成一个Waker。
二、整个故事的主线(最重要)
代码虽然多,但本质是这样一个闭环:
1. spawner 把「任务(含Future)」塞进任务队列
2. executor 从队列取出任务,poll 一下
3. Future 没做完 → 返回 Pending,但偷偷记下 waker
4. 未来某个线程完成任务 → 用 waker 唤醒 → 把任务重新塞回队列
5. executor 又取到它 → 再 poll → 这次返回 Ready,任务结束
下面我把这 5 步逐一展开。
三、逐个角色拆解
角色 A:TimerFuture(一个会"等 2 秒"的任务)
pub struct TimerFuture {
shared_state: Arc<Mutex<SharedState>>,
}
它的 new(duration) 做了什么:
- 创建一个共享状态
SharedState { completed: false, waker: None },用Arc<Mutex<>>包起来。 - 另外起一个新线程,这个线程里:
thread::sleep(duration)睡 2 秒。- 睡醒后,把
completed改成true。 - 如果里面存了
waker,就waker.wake()叫醒对方。
它的 poll 做了什么:
fn poll(self, cx) {
let mut shared_state = self.shared_state.lock().unwrap();
if shared_state.completed {
Poll::Ready(()) // 已经完成了
} else {
shared_state.waker = Some(cx.waker().clone()); // 记下"叫醒我"的方式
Poll::Pending // 还没完成,先撤
}
}
理解这个设计:poll 第一次被调用时,定时器肯定还没到 2 秒,所以返回 Pending。但它在返回前,把 cx.waker() 存了下来。等 2 秒后那个新线程睡醒,就拿着这个 waker 来"叫醒"它。
这里注意:
waker.wake()最终会触发什么?它会把Task重新塞回任务队列(下面角色 C 会解释)。
角色 B:Spawner(负责"接单")
struct Spawner {
task_sender: SyncSender<Arc<Task>>,
}
spawn(future) 做了什么:
- 把传进来的 future 用
.boxed()装箱成BoxFuture。 - 包成一个
Arc<Task>(Task 里面装着这个 future 和 sender)。 - 用
task_sender.send(task)把它塞进任务队列。
一句话:Spawner 就是"把新任务登记到队列"。
角色 C:Task(任务本体 + 唤醒桥梁)
struct Task {
future: Mutex<Option<BoxFuture<'static, ()>>>,
task_sender: SyncSender<Arc<Task>>,
}
这里有两个关键点:
① future 外面套了 Mutex<Option<...>>
Option:执行器 poll 时会把 futuretake()出来(变成None),poll 完如果还 Pending 再放回去(Some)。这样避免"边 poll 边改自己"的借用冲突。Mutex:因为 future 可能跨线程被唤醒,编译器要求线程安全。
② ArcWake 的实现 —— 这是整个唤醒机制的核心
impl ArcWake for Task {
fn wake_by_ref(arc_self: &Arc<Self>) {
let cloned = arc_self.clone();
arc_self.task_sender.send(cloned).expect("任务队列已满");
}
}
ArcWake 是 futures crate 提供的 trait,它让一个 Arc<Task> 可以被转换成 Waker。
wake_by_ref 的意思就是:"把我这个任务,重新塞回任务队列"。
于是整条唤醒链就通了:
新线程 waker.wake()
→ 内部调用 Task::wake_by_ref
→ 把这个 Task 的 Arc 重新 send 回队列
→ Executor 的 recv() 又能取到它了
角色 D:Executor(真正的"工人")
struct Executor {
ready_queue: Receiver<Arc<Task>>,
}
fn run(&self) {
while let Ok(task) = self.ready_queue.recv() {
let mut future_slot = task.future.lock().unwrap();
if let Some(mut future) = future_slot.take() {
let waker = waker_ref(&task); // ① 用 Task 造一个 waker
let context = &mut Context::from_waker(&*waker); // ② 包成 Context
if future.as_mut().poll(context).is_pending() { // ③ poll 一下
*future_slot = Some(future); // ④ 还没完,放回去
}
}
}
}
逐行解释 run():
while let Ok(task) = self.ready_queue.recv():不断从队列取任务。如果队列空了(所有 sender 都 drop 了),recv返回Err,循环结束。future_slot.take():把 future 从Option里"拿出来"(此时Option变成None)。waker_ref(&task):基于当前这个 Task 生成一个 waker。这个 waker 一旦被调用,就会把这个 Task 重新塞回队列(因为ArcWake就是这么实现的)。future.as_mut().poll(context):真正的"问一下"。- 如果返回
Pending:说明没做完,*future_slot = Some(future)把 future 放回 Task,等下次被唤醒再 poll。 - 如果返回
Ready:什么都不用做,future 已经消费掉了,Option保持None,这个 Task 就此消失。
四、把 main 那几行串起来走一遍(重点!)
现在看 main:
fn main() {
let (executor, spawner) = new_executor_and_spawner(); // 创建队列+两端
spawner.spawn(async { // 登记任务
println!("howdy!");
TimerFuture::newnew(2, 0).await;
println!("done!");
});
drop(spawner); // 关闭"接单端"
executor.run(); // 开始干活
}
时间线(完整执行过程):
阶段 1:准备工作
new_executor_and_spawner()创建了一个sync_channel,得到两个端点:task_sender(发送端)→ 放进Spawnerready_queue(接收端)→ 放进Executor
- 这个通道就是"任务队列"。
阶段 2:spawner.spawn(async { ... })
- 这个
async {}块会被编译器自动转换成一个匿名的 Future 结构体(一个状态机)。 - 这个状态机有 3 个状态点(对应 3 个
.await边界):- 还没开始
- 打印完
howdy!,正在等TimerFuture - 打印完
done!,结束
spawn把这个 future 装箱、包成Task、send进队列。- 此时只是"登记",一行代码都还没执行(没打印任何东西)。
阶段 3:drop(spawner)
- 把发送端
spawner丢掉。但注意:Task内部还clone了一份task_sender,所以通道还没真正关闭,队列还能继续运作。 - 这个
drop的意义是:当所有任务都完成后,最后一个task_sender也被销毁,通道彻底关闭,recv()才会返回Err,run()才会退出。
阶段 4:executor.run() 开始循环
第 1 轮循环:
recv()取到 Task(就是那个 async 块)。take()拿出 future。- 用 Task 造出 waker,包成 Context。
poll(context):- 状态机从"还没开始"进入执行,打印
howdy!。 - 走到
TimerFuture::new(...).await这里:- 先调用
TimerFuture::new(2s),这个函数马上启动一个新线程去睡觉,然后返回一个TimerFuture。 - 接着对这个
TimerFuture进行poll:completed还是false(才刚开始,2 秒没到)。- 于是把
cx.waker().clone()存进shared_state.waker。 - 返回
Poll::Pending。
- 先调用
- 因为内层的
TimerFuture是 Pending,整个 async 块也返回 Pending(.await会向上传播)。
- 状态机从"还没开始"进入执行,打印
is_pending()为真 →*future_slot = Some(future)把 future 放回 Task。
此时循环回到 recv(),队列是空的,run() 阻塞等待。
阶段 5:2 秒后,新线程醒来
thread::sleep(duration); // 睡够了
shared_state.completed = true; // 标记完成
if let Some(waker) = shared_state.waker.take() {
waker.wake() // 叫醒!
}
- 把
completed设为true。 - 取出之前存的 waker,调用
wake()。 wake()→Task::wake_by_ref→ 把这个 Task 重新send回队列。
阶段 6:run() 的第 2 轮循环
recv()又取到同一个 Task(被唤醒塞回来的)。take()拿出 future(这次状态机已经停在"等 TimerFuture"这个点)。- 重新造 waker、包 Context、
poll:- 状态机继续从
.await恢复,再次poll内层的TimerFuture。 - 这次
completed == true→ 返回Poll::Ready(())。 - 于是
.await完成,继续往下走,打印done!。 - async 块整体返回
Ready(())。
- 状态机继续从
is_pending()为假 → 不放回 future,Option保持None,Task 就此终结。
阶段 7:结束
- 最后一个 Task 完成了,它内部的
task_sender(Arc 里的)也被销毁。 - 通道彻底关闭,
recv()返回Err,run()的 while 循环退出,程序结束。
最终输出:
howdy!
(停顿 2 秒)
done!
五、核心思想总结(记住这个就懂了)
Future就是"分步执行的状态机",.await会把一个大的异步任务拆成多个暂停点。poll是"非阻塞地问一句":做完返回Ready,没做完返回Pending然后立刻返回,不占着线程傻等。Waker是"预约回拨":Future 暂停时,把cx.waker()存起来;等条件满足(如定时器到期),另一个线程调用wake(),把任务重新塞回队列。Executor是"循环工人":不停从队列取任务 → poll → 没做完就放回去等唤醒 → 做完就丢弃。- 整个异步模型 = "任务队列 + 轮询 + 唤醒",没有任何魔法,就是这三样东西在循环。
用一句话概括这段代码:
main只做了三件事:创建队列 → 登记任务 → 启动工人。工人poll任务时,任务说"我还没好,2 秒后叫醒我"(返回 Pending 并留下 waker)。2 秒后定时线程"叫醒"它(把任务塞回队列),工人再poll一次,任务这次说"我好了"(返回 Ready),结束。
为什么这么绕?—— 因为它是"隐式的协议",不是显式的调用
普通同步代码,逻辑是显式的、线性的:
A 调用 B,B 返回给 A,A 继续
你能一眼看到调用关系。但异步这套东西,真正干活的不是"函数调用",而是"约定俗成的协议":
| 隐式机制 | 它隐藏了什么 |
|---|---|
.await |
编译器偷偷把代码拆成状态机,你看到的"顺序代码"其实是"分段代码" |
waker.wake() |
表面是"叫醒",实际是"往队列塞东西",中间隔了好几层 |
Arc 的引用计数 |
表面是"一个指针",实际暗中决定"谁生谁死" |
recv() 的阻塞 |
表面是"取一个值",实际是"没活就睡,有活才醒" |
send / drop sender |
表面是"发个东西/丢个变量",实际暗中控制"通道什么时候关闭" |
核心难点就是:这些机制之间靠"副作用"互相牵动,而不是靠显式的调用线连在一起。 你要理解的是"谁在背后触发了谁",而不是"谁调用了谁"。
一个能反复用的方法:把"隐式链"翻译成"显式链"
当你看不懂一段异步代码时,问自己这几个固定问题,就能把链子串起来:
问题 1:这个动作的"真实副作用"是什么?
别管它表面叫什么,问它实际做了什么:
wake()→ 真实动作是"把 task 塞回队列"recv()→ 真实动作是"阻塞等活,有活才返回"clone()→ 真实动作是"引用计数 +1,给对象续命"drop(x)→ 真实动作是"引用计数 -1,可能触发销毁"
问题 2:这个副作用会"唤醒"谁?
每个副作用都会牵动另一个等待者:
send→ 唤醒阻塞在recv()上的执行器wake→ 唤醒阻塞在recv()上的执行器(本质同一个)drop sender→ 让recv()返回Err
问题 3:谁持有所有权?所有权转移时发生了什么?
Rust 异步里,所有权 = 生命周期。追踪"这个 Arc 现在在谁手里",就能知道"谁还活着"。
把这段代码的"隐形链"一次性摊开
我帮你把整段代码所有隐式的东西,串成一条显式的因果链。这是全文的"藏宝图":
【1】spawner.spawn(async { ... })
↓ 真实动作:future 装箱 → 包成 Arc<Task> → send 进队列
【2】drop(spawner)
↓ 真实动作:Spawner 的 sender 销毁(但 Task 里的 sender 还活着)
【3】executor.run() → recv()
↓ 真实动作:取到任务,开始阻塞式循环
【4】poll(async块)
↓ 真实动作:打印 howdy!,创建 TimerFuture,启动新线程去 sleep(2s)
【5】poll(TimerFuture)
↓ 真实动作:completed=false → 存下 waker → 返回 Pending
↓ waker 内部偷偷 clone 了一份 Arc<Task>(给 Task 续命)
【6】future 放回 Task,循环回到 recv()
↓ 真实动作:队列空 → 阻塞(不是退出!)
【7】2秒后,定时线程 sleep 结束
↓ 真实动作:completed=true → waker.wake()
↓ wake 内部:clone 一份 Arc<Task> → send 回队列
【8】recv() 被 send 唤醒,取到 Task
↓ 真实动作:再次 poll
【9】poll(TimerFuture)
↓ 真实动作:completed=true → 返回 Ready
【10】poll(async块) 继续 → 打印 done! → 返回 Ready
↓ 真实动作:future 不放回,Option=None
【11】循环体结束,局部变量 task drop
↓ 真实动作:最后一份 Arc 归零 → Task 销毁 → 里面 sender 销毁
【12】通道彻底关闭,recv() 返回 Err → run() 退出
每一步的右边,才是真正发生的事情。 左边那些 wake、recv、drop、poll 只是"表面名字"。
给你一个"降维打击"的视角
其实整个模型,剥掉所有 Rust 的包装,就只剩三件事在循环:
1. 有个队列(装任务)
2. 有个工人(循环从队列拿任务 poll)
3. 任务没做完时,留个"回拨方式"(waker),等条件满足时把自己塞回队列
Arc、Mutex、Sender、Pin、BoxFuture…… 全是为了让这三件事"跨线程、安全、不内存泄漏"而加的保险丝。核心逻辑就那三句话,其他都是"工程上的补丁"。
一个实操建议
下次再遇到这种绕的代码,别从头读到尾,而是倒着推:
- 先找"最终是谁让循环结束的"(这里:所有 sender drop)。
- 再往前追"谁持有 sender"(Task)。
- 再追"Task 什么时候销毁"(最后一份 Arc drop)。
- 再追"Arc 都有几份、在谁手里"(队列一份、waker 一份、局部变量一份)。
- 最后追"waker 是谁造的、什么时候被调用"(
waker_ref(&task)造的,定时线程调用)。
从"结果"倒推"原因",比顺着读更容易把隐式链串起来,因为隐式的东西往往是在"结果处"才显现出意义的。
{补充:async 块会被编译器编译成一个状态机(impl Future 的结构体)。每个 .await 是一个暂停点,把代码切成若干段;跨越暂停点的局部变量会被存进状态机结构体里,状态机用一个整数记录"当前执行到哪一段"。当 Future 因 Pin 被"钉住"(防止自引用悬垂)后返回 Pending 时,它会把 Waker 存下来。等到条件满足,wake() 把任务(Arc<Task>)重新塞回队列,执行器被唤醒后再次 poll,状态机就从上次的暂停点继续往下走,直到返回 Ready。)