迷你的异步运行时

你提到的"引入了许多外部方法",其实就是这几个标准库的东西,我用大白话解释:

1. Future 特征(trait)

trait Future {
    type Output;
    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}

2. Waker 和 Context

3. Pin<&mut Self>

4. Arc / Mutex

5. sync_channel(同步通道)

6. BoxFuture / .boxed() / waker_ref / ArcWake


二、整个故事的主线(最重要)

代码虽然多,但本质是这样一个闭环:

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) 做了什么:

  1. 创建一个共享状态 SharedState { completed: false, waker: None },用 Arc<Mutex<>> 包起来。
  2. 另外起一个新线程,这个线程里:
    • 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) 做了什么:

  1. 把传进来的 future 用 .boxed() 装箱成 BoxFuture。
  2. 包成一个 Arc<Task>(Task 里面装着这个 future 和 sender)。
  3. 用 task_sender.send(task) 把它塞进任务队列。

一句话:Spawner 就是"把新任务登记到队列"。


角色 C:Task(任务本体 + 唤醒桥梁)

struct Task {
    future: Mutex<Option<BoxFuture<'static, ()>>>,
    task_sender: SyncSender<Arc<Task>>,
}

这里有两个关键点:

① future 外面套了 Mutex<Option<...>>

② 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():

  1. while let Ok(task) = self.ready_queue.recv():不断从队列取任务。如果队列空了(所有 sender 都 drop 了),recv 返回 Err,循环结束。
  2. future_slot.take():把 future 从 Option 里"拿出来"(此时 Option 变成 None)。
  3. waker_ref(&task):基于当前这个 Task 生成一个 waker。这个 waker 一旦被调用,就会把这个 Task 重新塞回队列(因为 ArcWake 就是这么实现的)。
  4. future.as_mut().poll(context):真正的"问一下"。
  5. 如果返回 Pending:说明没做完,*future_slot = Some(future) 把 future 放回 Task,等下次被唤醒再 poll。
  6. 如果返回 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:准备工作

阶段 2:spawner.spawn(async { ... })

阶段 3:drop(spawner)

阶段 4:executor.run() 开始循环

第 1 轮循环:

  1. recv() 取到 Task(就是那个 async 块)。
  2. take() 拿出 future。
  3. 用 Task 造出 waker,包成 Context。
  4. 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 会向上传播)。
  5. 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()                   // 叫醒!
}

阶段 6:run() 的第 2 轮循环

  1. recv() 又取到同一个 Task(被唤醒塞回来的)。
  2. take() 拿出 future(这次状态机已经停在"等 TimerFuture"这个点)。
  3. 重新造 waker、包 Context、poll:
    • 状态机继续从 .await 恢复,再次 poll 内层的 TimerFuture。
    • 这次 completed == true → 返回 Poll::Ready(())。
    • 于是 .await 完成,继续往下走,打印 done!。
    • async 块整体返回 Ready(())。
  4. is_pending() 为假 → 不放回 future,Option 保持 None,Task 就此终结。

阶段 7:结束

最终输出:

howdy!
(停顿 2 秒)
done!

五、核心思想总结(记住这个就懂了)

  1. Future 就是"分步执行的状态机",.await 会把一个大的异步任务拆成多个暂停点。
  2. poll 是"非阻塞地问一句":做完返回 Ready,没做完返回 Pending 然后立刻返回,不占着线程傻等。
  3. Waker 是"预约回拨":Future 暂停时,把 cx.waker() 存起来;等条件满足(如定时器到期),另一个线程调用 wake(),把任务重新塞回队列。
  4. Executor 是"循环工人":不停从队列取任务 → poll → 没做完就放回去等唤醒 → 做完就丢弃。
  5. 整个异步模型 = "任务队列 + 轮询 + 唤醒",没有任何魔法,就是这三样东西在循环。

用一句话概括这段代码:

main 只做了三件事:创建队列 → 登记任务 → 启动工人。工人 poll 任务时,任务说"我还没好,2 秒后叫醒我"(返回 Pending 并留下 waker)。2 秒后定时线程"叫醒"它(把任务塞回队列),工人再 poll 一次,任务这次说"我好了"(返回 Ready),结束。


为什么这么绕?—— 因为它是"隐式的协议",不是显式的调用

普通同步代码,逻辑是显式的、线性的:

A 调用 B,B 返回给 A,A 继续

你能一眼看到调用关系。但异步这套东西,真正干活的不是"函数调用",而是"约定俗成的协议":

隐式机制 它隐藏了什么
.await 编译器偷偷把代码拆成状态机,你看到的"顺序代码"其实是"分段代码"
waker.wake() 表面是"叫醒",实际是"往队列塞东西",中间隔了好几层
Arc 的引用计数 表面是"一个指针",实际暗中决定"谁生谁死"
recv() 的阻塞 表面是"取一个值",实际是"没活就睡,有活才醒"
send / drop sender 表面是"发个东西/丢个变量",实际暗中控制"通道什么时候关闭"

核心难点就是:这些机制之间靠"副作用"互相牵动,而不是靠显式的调用线连在一起。 你要理解的是"谁在背后触发了谁",而不是"谁调用了谁"。


一个能反复用的方法:把"隐式链"翻译成"显式链"

当你看不懂一段异步代码时,问自己这几个固定问题,就能把链子串起来:

问题 1:这个动作的"真实副作用"是什么?

别管它表面叫什么,问它实际做了什么:

问题 2:这个副作用会"唤醒"谁?

每个副作用都会牵动另一个等待者:

问题 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…… 全是为了让这三件事"跨线程、安全、不内存泄漏"而加的保险丝。核心逻辑就那三句话,其他都是"工程上的补丁"。


一个实操建议

下次再遇到这种绕的代码,别从头读到尾,而是倒着推:

  1. 先找"最终是谁让循环结束的"(这里:所有 sender drop)。
  2. 再往前追"谁持有 sender"(Task)。
  3. 再追"Task 什么时候销毁"(最后一份 Arc drop)。
  4. 再追"Arc 都有几份、在谁手里"(队列一份、waker 一份、局部变量一份)。
  5. 最后追"waker 是谁造的、什么时候被调用"(waker_ref(&task) 造的,定时线程调用)。

从"结果"倒推"原因",比顺着读更容易把隐式链串起来,因为隐式的东西往往是在"结果处"才显现出意义的。


{补充:async 块会被编译器编译成一个状态机(impl Future 的结构体)。每个 .await 是一个暂停点,把代码切成若干段;跨越暂停点的局部变量会被存进状态机结构体里,状态机用一个整数记录"当前执行到哪一段"。当 Future 因 Pin 被"钉住"(防止自引用悬垂)后返回 Pending 时,它会把 Waker 存下来。等到条件满足,wake() 把任务(Arc<Task>)重新塞回队列,执行器被唤醒后再次 poll,状态机就从上次的暂停点继续往下走,直到返回 Ready。)