Exploring-Async-Basics-with-Rust

Whenever you have eliminated the impossible, whatever remains, however improbable, must be the truth.

— Arthur Conan Doyle


photo by Samuel Field(https://unsplash.com/@s_giovanni?utm_source=templater_proxy&utm_medium=referral) on Unsplash

第 1 章:并发 vs 并行

1.1 核心定义

概念 一句话定义 关注点
并发 (Concurrency) 同一时间处理(dealing with)许多事情 资源的效用与效率
并行 (Parallelism) 同一时间(doing)许多事情 增加资源投入
多任务 (Multitasking) 同一时间使多个任务取得进展 (progress) 实现方式:并发 或 并行

经典比喻:并发是更聪明地工作;并行是投入更多资源到问题上。

1
2
3
4
5
6
7
8
9
并发(单核交替推进,任务"可中断"):
Task A: ████----████------████
Task B: ----████----██████----
时间轴: ──────────────────────→ (虚线 = 暂停等待)

并行(多核同时执行):
Core 1 / Task A: ████████████████
Core 2 / Task B: ████████████████
时间轴: ────────────────→

1.2 为什么需要并发(两大场景)

1
2
3
4
5
6
7
8
9
10
// 场景 1:I/O 等待 —— 等待外部事件时不能干等
// 查询数据库时 CPU 闲置,应切去跑别的任务
db_query(sql, |result| { /* 数据回来再继续 */ });

// 场景 2:公平性 —— 防止一个任务饿死其他任务(典型:UI 响应)
// 单核上每 16ms 打断一次重计算,保证 UI 以 60Hz 刷新
loop {
heavy_task.run_for(16.ms);
ui.refresh(); // 每 16ms 让 UI "进展" 一次
}

1.3 关键洞察:参照系(Frame of Reference)

"同步执行"是一场错觉——它只在"你写的代码"这个参照系中成立:

graph TD
    subgraph 三层参照系
        A["程序员视角<br/>自己的代码按顺序执行<br/>'看起来是同步的'"]
        B["OS 视角<br/>抢占式调度<br/>随时暂停/恢复你的线程"]
        C["CPU 视角<br/>随时被硬件中断<br/>流水线+乱序执行"]
    end
    A -->|"其实被"| B -->|"其实被"| C

谈论并发时若不先声明参照系,很快就会陷入混乱——这是理解后文一切的基础。


第 2 章:异步的历史

2.1 演进时间线

graph LR
    A["单 CPU<br/>顺序执行"] --> B["非抢占式多任务<br/>(协作式, DOS/Win95)"]
    B --> C["抢占式多任务<br/>(OS 负责调度, 现代 OS)"]
    C --> D["超线程<br/>1 核模拟 2 逻辑核"]
    D --> E["多核处理器<br/>+ 每核超线程"]

2.2 非抢占式 vs 抢占式多任务

1
2
3
4
5
6
7
8
9
10
11
12
13
// ── 非抢占式(协作式):程序员必须主动交还控制权 ──
fn window_message_loop() {
loop {
handle_one_event();
yield_to_os(); // 忘记写这行 → 整个系统卡死!
} // (Win95 时代拖拽卡死的窗口能画满全屏)
}

// ── 抢占式:OS 强行收走 CPU,程序员无感知 ──
fn my_code() {
let x = 1 + 1; // OS 可能在任意两条指令间暂停本线程
let y = x * 2; // 去运行鼠标更新/其他进程,然后再切回来
}
  • 非抢占式:调度责任在程序员,一个 bug 拖垮整个系统
  • 抢占式:OS 每秒多次上下文切换,保证 UI/后台任务/I/O 都拿到时间片(现代 OS 的标准设计)

2.3 超线程(Hyperthreading)

1
2
3
4
一个物理核内有多个运算单元(如 ALU):
线程 1 使用 ALU 时,线程 2 可利用闲置的浮点单元/逻辑单元
→ 1 个物理核 = 2 个逻辑核(如 6 核 12 线程)
→ 额外性能约 +30%(取决于工作负载,自 1990s 持续改进)

2.4 "你的代码有多同步?"——三个视角再回顾

  • 程序/线程视角:按编写顺序执行 ✔(同步的错觉)
  • OS 视角:可能中断、暂停、恢复你的代码
  • CPU 视角:流水线(当前指令执行时预取下一条)、分支预测、乱序重排指令(不告知程序员/OS)——所以"A 发生在 B 之前"在硬件层并无保证,这正是同步原语(mutex/atomic)存在的原因

第 3 章:OS 与 CPU

3.1 OS"假装"同步

OS 自 1990s 以来就"假装"以同步方式执行程序。所谓同步代码只是对程序员看似同步;OS 采用抢占式调度,你无法保证代码不被中断地逐条执行。

3.2 系统调用(syscall)——与 OS 通信的唯一正道

你无权直接操作硬件(如网卡),必须通过 syscall 请求 OS 代劳。

三个抽象层次实现同一个功能(向 stdout 输出):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// ── 层次 1:内联汇编(直接给 CPU 下指令)──
// Linux: write 是 1 号系统调用, 1 也是 stdout 的 fd
unsafe {
llvm_asm!("
mov $$1, %rax # syscall 号: write
mov $$1, %rdi # fd: stdout
mov $0, %rsi # 缓冲区地址
mov $1, %rdx # 长度
syscall # 进入内核
" : : "r"(msg_ptr), "r"(len) : "rax","rdi","rsi","rdx");
}
// macOS 几乎相同, 仅系统调用号不同: 0x2000004
// Windows: 无稳定调用号保证(逆向工程表), 绝不可用此方式!

// ── 层次 2:链接 OS 提供的 C 库/WinAPI ──
#[cfg(not(target_os = "windows"))]
#[link(name = "c")]
extern "C" { // UNIX: libc, C 调用约定
fn write(fd: u32, buf: *const u8, count: usize) -> i32;
}

#[cfg(target_os = "windows")]
#[link(name = "kernel32")]
extern "stdcall" { // Windows: stdlib 调用约定, 更复杂
fn GetStdHandle(nStdHandle: i32) -> i32; // -11 = stdout
fn WriteConsoleW(h: i32, buf: *const u16, ...); // 需 utf-16!
}

// ── 层次 3:标准库(最高抽象)──
println!("Hello world from Stdlib");

要点:

  • syscall 指令比早期 int 0x80 软中断快,借助 VDSO(附属于进程的内存页)避免上下文切换
  • Windows 的 WriteConsoleW 需要 utf-16,故需 message.encode_utf16() 转换;W=Unicode,A=ANSI
  • FFI 调用必须 unsafe;Linux/macOS API 简单相似,Windows 需要更多数据结构与代码

3.3 跨平台抽象的隐藏复杂性

代码只跑 Linux/macOS 很省事,一旦跨平台代码量爆炸。社区方案:libc crate(封装函数+常量)、miolibuv。边缘情况(如"所有合法 utf-8 都能转 utf-16 正确显示吗?")是额外的难题。

3.4 CPU 与 OS 的暗中合作(安全机制)

一个"为什么 CPU 知道不允许解引用 99999999999999"的实验引出整套硬件级安全机制:

graph TD
    A["代码解引用指针"] --> B["MMU 查页表<br/>虚拟地址 → 物理地址"]
    B -->|"找到映射"| C["正常读取数据 ✔"]
    B -->|"页表中无映射"| D["触发 Page Fault 异常"]
    D --> E["CPU 查中断描述符表 IDT<br/>(OS 启动时注册)"]
    E --> F["跳转到 OS 的处理程序<br/>输出: segmentation fault"]

    G["Ring 0 内核态"] -.->|"可改页表、访问设备"| B
    H["Ring 3 用户态"] -.->|"受限访问, 违规即异常"| B
  • 虚拟内存 + 页表:每个进程有独立页表,CPU 特殊寄存器指向它
  • 特权级 Ring 0(内核)/ Ring 3(用户):用户态代码试图改页表 → 异常 → 跳转到 OS 处理程序。这就是除系统调用外你别无选择与内核/硬件交互的原因

第 4 章:中断 | 固件 | I/O

4.1 从网卡读取数据的全流程(简化)

graph TD
    A["1. 我们的代码<br/>syscall 注册 socket<br/>(得 fd / handle)"] --> B["2. 向 OS 注册<br/>感兴趣的事件"]
    B --> C["3. 网卡固件<br/>专用微控制器检测到<br/>数据到达"]
    C --> D["4. 硬件中断<br/>IRQ 电信号发给 CPU"]
    D --> E["5. CPU 查 IDT<br/>跳到中断处理程序"]
    E --> F["6. DMA<br/>网卡直接把数据<br/>写入内存缓冲区"]
    F --> G["7. 驱动程序<br/>通知内核数据就绪"]
    G --> H["8. OS 按注册方式通知我们"]

4.2 第 2 步"注册事件"的三种姿势

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// A) 阻塞式: OS 挂起我的线程, 直到事件发生
// → 简单, 但线程被挂起干不了别的
blocking_read(fd);

// B) 非阻塞式: 立即返回句柄, 自己负责轮询检查
// → 不阻塞, 但要忙轮询(浪费 CPU)
match nonblocking_poll(fd) {
Ready => read(fd),
NotReady => /* 再等等, 周期性再查 */
}

// C) 事件队列: 一次订阅多个事件, 轮询队列时阻塞
// → epoll/kqueue/IOCP 的模型, 兼顾效率与灵活
event_queue.register(fd, Read);
for event in event_queue.wait() { // 阻塞直到任一事件发生
read(event.fd);
}

4.3 硬件中断 / 软件中断 / 固件

  • 硬件中断:IRQ 线上的电信号,随时打断 CPU → 存寄存器状态 → 查 IDT 跳转处理程序
  • 软件中断:由指令主动触发(如旧式 int 0x80 系统调用)
  • IDT 位于主内存的固定位置,CPU 只在寄存器存指向它的指针(网卡类中断的处理程序通常由驱动程序注册)
  • 固件:系统中"隐藏的小 CPU"无处不在(网卡、甚至 CPU 自身)。网卡固件已在轮询数据,我们若再让 CPU 忙轮询就是重复劳动——并发即效率的又一体现

第 5 章:处理 I/O 的三大策略

策略对比总表

策略1: OS 线程 策略2: 绿色线程 策略3: OS 事件队列(epoll/kqueue/IOCP)
代表 传统服务器 Go Node / Rust async
优点 简单易写;性能尚可;免费获得并行 用法像 OS 线程;可控调度/优先级 接近最佳资源利用率;最大灵活性
缺点 栈大(海量任务耗尽内存);syscall 开销大;OS 调度不可控 需要运行时,重复 OS 已做的事;实现难 各 OS 差异大;复杂;只解决"等待",仍需暂停任务的方案
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 策略 1:一任务一线程(简单粗暴)
for conn in incoming {
thread::spawn(move || handle(conn)); // 10k 连接 = 10k 线程 ≈ 内存爆炸
}

// 策略 2:绿色线程(Go 风格,运行时自己调度轻量线程)
go!(|| handle(conn)); // 有运行时, 自己管栈和调度

// 策略 3:OS 事件队列(本书主角)
epoll.register(conn, READABLE, token);
loop { // 单线程管理所有连接
for ev in epoll.wait() { // 阻塞直到有事件
dispatch(ev); // 事件到了再处理
}
}

关键结论:策略 3 只解决了"何时知道就绪",如何暂停/恢复任务仍需上层方案——

  • Node 的答案:回调 (callbacks)
  • Rust 的答案:Futures 状态机(暂停点即状态)

Node 运行时 = 策略 1 + 策略 3 组合,但尽量强制 I/O 走策略 3——这是 Node 擅长海量并发连接的原因。


第 6 章:Epoll | Kqueue | IOCP

6.1 谁在用它们

1
2
3
Node  → libuv   (跨平台异步 I/O 库, Julia/Pyuv 也用)
Rust → mio (tokio 的 OS 事件队列后端; tokio之于Actix Web...)
本书 → minimio (作者手写的玩具版 mio)

6.2 两种模型的根本差异

sequenceDiagram
    participant P as epoll/kqueue (就绪型)
    Note over P: 告诉你"可以读了", 数据你自己取
    P->>P: epoll_create()/kqueue() 建队列
    P->>P: 注册 fd 的 Read 兴趣
    P->>P: epoll_wait()/kevent() 阻塞等待
    P-->>P: 返回"socket N 可读"
    P->>P: 自己调 read(N)

    participant C as IOCP (完成型)
    Note over C: 告诉你"已经读好了", 数据已在缓冲区
    C->>C: CreateIoCompletionPort() 建队列
    C->>C: 注册 Read + 借出缓冲区给 OS
    C->>C: GetQueuedCompletionStatusEx() 阻塞
    C-->>C: 返回"数据已读入你的缓冲区"
Epoll (Linux) Kqueue (macOS/BSD) IOCP (Windows)
模型 就绪型 readiness 就绪型 readiness 完成型 completion
通知含义 "可以执行操作了" 同左 "操作已完成,数据已进缓冲区"
缓冲区 自己管理 自己管理 借给 OS(等待期间不可动,借用检查器大显身手)
特点 大量事件下高效 概念类似但更抽象通用 官方文档好

跨平台库设计经验:让"就绪型"表现得像"完成型"更容易 → 先按 IOCP(完成型)设计,再适配 epoll/kqueue


第 7 章:例子

7.1 目标代码(JS 风格的 Rust)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 把这个函数想象成你写的 JavaScript 程序
fn javascript() {
Fs::read("test.txt", |result| { // 回调式异步
let text = result.into_string().unwrap();
Crypto::encrypt(text.len(), |result| { // 回调嵌套(类似 JS)
print(result.into_int().unwrap());
});
});
set_timeout(0, |_| print("Immediate1")); // 计时器
Http::http_get_slow("www.google.com", 2000, |result| {
print_content(result.into_string().unwrap(), "web call");
});
}

fn main() {
let rt = Runtime::new();
rt.run(javascript); // 运行时驱动一切
}

Js 枚举模拟 JS 的动态类型(Undefined/String/Int),因为 Rust 是静态类型语言。

7.2 什么是 Node(破除误解)

graph TD
    subgraph Node 架构
        V["V8 引擎<br/>JIT 编译执行 JS<br/>(不能 I/O)"]
        R["Node 运行时"]
        EL["事件循环<br/>(用户代码只跑在这一个线程!)"]
        TP["线程池 x4<br/>(默认)"]
        EQ["libuv<br/>epoll/kqueue/IOCP 事件队列"]
    end
    JS["你的 JS 代码"] --> V --> R
    R --> EL
    EL -->|"I/O 密集"| EQ
    EL -->|"CPU 密集 / epoll 处理不了的 I/O<br/>(如文件读取)"| TP
  • JS 本身没有事件循环——浏览器或 Node 提供运行时才有
  • Node 是多线程的(线程池 + epoll 线程),但你的代码只跑在单个主线程上——"别阻塞事件循环"指的就是它
  • 受 I/O 限制 → 事件队列;受 CPU 限制 → 线程池(大多数 C++ 扩展也在此执行)

7.3 实现计划

需要:两个事件队列(线程池 + minimio)|一个运行时(存回调、发任务、注册事件、轮询、处理计时器)|三个模块(Fs/Crypto/Http)|辅助函数(print/print_content/current


第 8 章:实现自己的运行时

8.1 运行时的状态(Runtime struct 逐字段)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
pub struct Runtime {
available_threads: Vec<usize>, // 线程池空闲线程 id
callbacks_to_run: Vec<(usize, Js)>, // 本 tick 待运行的 (回调id, 数据)
callback_queue: HashMap<usize, Box<dyn FnOnce(Js)>>, // 所有已注册回调
epoll_pending_events: usize, // epoll 待处理计数(仅打印用)
epoll_registrator: minimio::Registrator, // 向 OS 注册事件
epoll_thread: thread::JoinHandle<()>, // epoll 线程句柄
epoll_timeout: Arc<Mutex<Option<i32>>>, // None=无限等待 Some(n)=n ms
event_reciever: Receiver<PollEvent>, // 线程池+epoll 共同的事件通道!
identity_token: usize, // 回调 id 生成器
pending_events: usize, // 归零 => 程序结束
thread_pool: Vec<NodeThread>, // 线程池句柄
timers: BTreeMap<Instant, usize>, // 到期时间 -> 回调id (有序!)
timers_to_remove: Vec<Instant>, // 复用缓冲,避免循环内分配
}

支撑类型:

1
2
3
4
5
6
7
8
9
10
11
12
13
struct Task {                              // 发给线程池的工作单元
task: Box<dyn Fn() -> Js + Send>, // 真正的活(Send 才能跨线程)
callback_id: usize, // 干完活回调谁
kind: ThreadPoolTaskKind, // FileRead | Encrypt | Close
}

enum PollEvent { // 主循环统一收到的事件
Threadpool((thread_id, callback_id, Js)), // 线程池完成(带数据)
Epoll(event_id), // epoll 事件就绪
Timeout, // epoll 等待超时(该查计时器了)
}

static mut RUNTIME: *mut Runtime = std::ptr::null_mut(); // 全局运行时指针(unsafe hack)

8.2 主循环(仿 Node 的 6 阶段)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
pub fn run(mut self, f: impl Fn()) {
unsafe { RUNTIME = &mut self }; // 全局指针 hack(单线程访问,安全)
f(); // 先同步跑一遍"JS程序"→ 注册一堆事件

while self.pending_events > 0 { // 没有待处理事件 => 结束
// 1. TIMERS: 到期计时器的回调入队
self.process_expired_timers();

// 2. CALLBACKS: 跑掉所有排队的回调
self.run_callbacks();

// 3. IDLE/PREPARE (Node 内部用, 略)

// 4. POLL: 计算最近一个计时器的超时, 写入共享变量
let next_timeout = self.get_next_timer();
*self.epoll_timeout.lock().unwrap() = next_timeout;
drop(lock); // ★ 必须先放锁再等待,否则死锁!

// 阻塞等待: 线程池完成 / epoll 就绪 / 超时, 三者任一
if let Ok(event) = self.event_reciever.recv() {
match event {
Timeout => (), // 转回去查计时器
Threadpool((tid, cb_id, data)) => self.process_threadpool_events(..),
Epoll(event_id) => self.process_epoll_events(..),
}
}
self.run_callbacks(); // 事件转化出的回调立刻跑

// 5. CHECK (setImmediate 钩子, 略) 6. CLOSE CALLBACKS (略)
}
// 清理: 发 Close 任务给每个池线程并 join, 关闭 epoll 循环并 join
}
graph TD
    S["run(f) 先同步执行用户代码<br/>注册 timers/线程池任务/epoll 事件"] --> L{"pending_events > 0 ?"}
    L -->|"否"| Z["清理并退出"]
    L -->|"是"| T1["1. Timers<br/>到期回调入队"]
    T1 --> T2["2. Callbacks<br/>执行排队的回调"]
    T2 --> T4["4. Poll<br/>设 timeout=最近计时器<br/>阻塞 recv 事件通道"]
    T4 --> E{"收到什么事件?"}
    E -->|"Timeout"| T1
    E -->|"Threadpool 完成"| T2
    E -->|"Epoll 就绪"| T2

8.3 线程池的搭建(Runtime::new 上半场)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 1 条主通道: 线程池+epoll线程 → 主循环
let (event_sender, event_receiver) = channel::<PollEvent>();

for i in 0..4 { // Node 默认 4 线程
let (evt_sender, evt_receiver) = channel::<Task>(); // 每线程独立通道
let event_sender = event_sender.clone();

let handle = thread::Builder::new().name(format!("pool{}", i))
.spawn(move || {
while let Ok(task) = evt_receiver.recv() { // recv 会 park 线程(零消耗)
if let ThreadPoolTaskKind::Close = task.kind { break; }
let res = (task.task)(); // 干活,返回 Js
event_sender.send(PollEvent::Threadpool((i, task.callback_id, res)));
} // ↑ 干完通知主循环
}).unwrap();
threads.push(NodeThread { handle, sender: evt_sender });
}

8.4 epoll 线程(Runtime::new 下半场)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
let mut poll = minimio::Poll::new().unwrap();     // syscall 向 OS 要 epoll/kqueue/IOCP 句柄
let registrator = poll.registrator(); // 注册器与 Poll 分离:
// 主线程持有 Registrator,
// Poll 发给 epoll 线程
// (共享 AtomicBool 判断队列是否存活)
let epoll_timeout = Arc::new(Mutex::new(None)); // 与主线程共享的超时值

thread::Builder::new().name("epoll".into()).spawn(move || {
let mut events = minimio::Events::with_capacity(1024); // 只分配一次!
loop {
let timeout = *epoll_timeout_clone.lock().unwrap();
drop(lock); // ★ 同样: 先放锁再 poll

match poll.poll(&mut events, timeout) { // OS 级阻塞等待(线程被 park)
Ok(v) if v > 0 => { /* 就绪事件逐个 send PollEvent::Epoll(id) */ }
Ok(0) => { /* 超时/伪唤醒: send PollEvent::Timeout */ }
Err(Interrupted) => break, // 收到关闭信号
Err(e) => panic!(e),
}
}
});

两个 drop(lock)全书最重要的并发细节poll/recv 会阻塞线程,若持有锁去等待,另一端将永远无法写入 → 死锁

8.5 Timers:BTreeMap 的妙用

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
fn process_expired_timers(&mut self) {
// BTreeMap 按到期时间有序 → range 一刀切出所有已到期者
self.timers.range(..=Instant::now())
.for_each(|(k, _)| timers_to_remove.push(*k)); // 先收集(借用检查)

while let Some(key) = self.timers_to_remove.pop() { // 再删除
let cb_id = self.timers.remove(&key).unwrap();
self.callbacks_to_run.push((cb_id, Js::Undefined)); // 只入队不执行
}
}

fn get_next_timer(&self) -> Option<i32> {
self.timers.iter().nth(0).map(|(&instant, _)| // 首个 = 最近到期
(instant - Instant::now()).as_millis() as i32)
}

BTreeMap 而非 HashMap/BST:键天然有序;相比 BST 每节点存多个值(小 Vec),缓存友好

8.6 Callbacks:Box<dyn FnOnce(Js)>

1
2
3
4
5
6
7
fn run_callbacks(&mut self) {
while let Some((cb_id, data)) = self.callbacks_to_run.pop() {
let cb = self.callback_queue.remove(&cb_id).unwrap();
cb(data); // ★ 回调里的长代码会阻塞整个事件循环!
self.pending_events -= 1; // 这就是"别阻塞事件循环"的本质
}
}

FnOnce:闭包捕获并消耗环境变量,只能调用一次——正是回调语义,还顺便借 RAII 清理资源。trait 无固定大小 → 必须 Box 上堆。

8.7 三个注册 API(模块与运行时的接口)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
// ① 注册 epoll 事件 (Http 模块用)
fn register_event_epoll(&mut self, token: usize, cb: impl FnOnce(Js)) {
self.add_callback(token, cb); // token 即回调 id (1:1 映射)
self.pending_events += 1;
}

// ② 注册线程池任务 (Fs/Crypto 模块用)
fn register_event_threadpool(&mut self, task: impl Fn()->Js+Send,
kind: ThreadPoolTaskKind, cb: impl FnOnce(Js)) {
let cb_id = self.generate_cb_identity();
self.add_callback(cb_id, cb);
let available = self.get_available_thread(); // 无空闲线程直接 panic(书中捷径)
self.thread_pool[available].sender.send(Task { task: Box::new(task), .. });
self.pending_events += 1;
}

// ③ set_timeout (计时器)
fn set_timeout(&mut self, ms: u64, cb: impl Fn(Js)) {
let cb_id = self.generate_cb_identity();
self.add_callback(cb_id, cb);
self.timers.insert(Instant::now() + Duration::from_millis(ms), cb_id);
self.pending_events += 1;
}

事件回流的处理:

1
2
3
4
5
6
7
fn process_threadpool_events(&mut self, thread_id, callback_id, data: Js) {
self.callbacks_to_run.push((callback_id, data)); // 结果数据直接带给回调
self.available_threads.push(thread_id); // ★ 线程归还池子
}
fn process_epoll_events(&mut self, event_id: usize) {
self.callbacks_to_run.push((event_id, Js::Undefined)); // 不带数据! 只报"就绪"
}

安全扩展:完成型模型(IOCP)要为每个 Read/Write 借出缓冲区,恶意客户端可借注册大量不发生的事件耗尽内存——生产实现需设"高水位线"限制待处理事件数。

8.8 基础设施杂项

1
2
3
4
5
fn generate_identity(&mut self) -> usize {
self.identity_token = self.identity_token.wrapping_add(1); // 溢出回绕
self.identity_token
}
// generate_cb_identity: 撞 id 就循环再生成(64 位上每纳秒一个 id 也要 585 年才绕完一圈)

第 9 章:模块(相当于 Node 的 C++ 扩展)

9.0 全局 RUNTIME 的 unsafe 及为何安全

1
2
3
4
5
6
pub fn set_timeout(ms: u64, cb: impl Fn(Js) + 'static) {
let rt = unsafe { &mut *(RUNTIME as *mut Runtime) }; // 模拟 JS 的全局函数
rt.set_timeout(ms, cb);
}
// 安全性依据: ① RUNTIME 只在一个线程访问
// ② 模块只会在 Runtime.run() 内被调用, 此时指针必有效

9.1 Fs 模块 + 彩蛋:为什么文件 I/O 走线程池而不是 epoll?

1
2
3
4
5
6
7
8
9
10
11
impl Fs {
fn read(path: &'static str, cb: impl Fn(Js)) {
let work = move || {
thread::sleep(Duration::from_secs(1)); // 模拟大文件,便于观察
let mut buffer = String::new();
fs::File::open(&path).unwrap().read_to_string(&mut buffer).unwrap();
Js::String(buffer) // 结果作为回调的输入
};
unsafe { &mut *RUNTIME }.register_event_threadpool(work, FileRead, cb);
}
}

文件 I/O 走线程池的原因:

  1. OS 大量缓存文件,读通常立即就绪——注册事件+等通知反而比直接读更慢
  2. 但"从 OS 缓存拷进你的缓冲区"仍耗时,会阻塞主循环 → 丢给线程池
  3. Linux/macOS 的就绪型模型对异步文件支持差;Windows(IOCP)/新版 Linux(io_uring) 才有好的完成型支持
1
2
线程池方案:  代码简单 ✔ | 性能够好 ✔ | 多数场景(如 web 服务器)几乎无损失
异步文件I/O: 复杂度高 ✘ | API 贫乏且平台差异大 ✘ | 实际收益甚微

9.2 Crypto 模块(CPU 密集型代表)

1
2
3
4
5
6
7
8
9
10
11
impl Crypto {
fn encrypt(n: usize, cb: impl Fn(Js)) {
let work = move || {
fn fib(n: usize) -> usize { // 递归斐波那契:
match n { 0 => 0, 1 => 1, _ => fib(n-1) + fib(n-2) }
} // 经典的烧 CPU workload
Js::Int(fib(n))
};
unsafe { &mut *RUNTIME }.register_event_threadpool(work, Encrypt, cb);
}
}

9.3 Http 模块(epoll 路线的代表)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
impl Http {
fn http_get_slow(url: &str, delay_ms: u32, cb: impl Fn(Js)) {
let rt = unsafe { &mut *RUNTIME };
// minimio::TcpStream 而非 std 的: 为了跨平台抽象(epoll就绪/IOCP完成两种模型)
let mut stream = minimio::TcpStream::connect("slowwly...:80").unwrap();
stream.write_all(format!("GET /delay/{} ... HTTP/1.1\r\n...", delay_ms).as_bytes());
// ↑ 书中捷径: 写是阻塞的(真实实现写也应异步)

let token = rt.generate_cb_identity(); // token == callback_id
rt.epoll_registrator.register(&mut stream, token,
minimio::Interests::READABLE).unwrap();
// ↑ 需要 &mut stream 是为了 IOCP: 借出缓冲区期间借用检查器保证没人碰它

let wrapped = move |_| { // 包装: 就绪后再真正 read
let mut buffer = String::new();
stream.read_to_string(&mut buffer);
cb(Js::String(buffer));
};
rt.register_event_epoll(token, wrapped);
}
}

伪唤醒 (spurious wakeup):OS 可能"不确定事件是否发生"时依然唤醒线程。契约要求程序员重注册事件再检查,而非假设数据已到。真实实现应在读到 WouldBlock 时重新注册读事件。


第 10~11 章:组装与最终代码

10.1 入口与全局架构图

1
2
3
4
fn main() {
let rt = Runtime::new(); // 建 4 线程池 + epoll 线程 + 共享通道
rt.run(javascript); // 跑"JS程序", 进入事件循环直到 pending_events==0
}
graph TD
    subgraph 主线程 main
        M["Runtime::run 事件循环<br/>timers→callbacks→poll→..."]
        CB["callback_queue / callbacks_to_run"]
        T["timers (BTreeMap)"]
    end

    subgraph 工作者
        P0["pool0"]
        P1["pool1~3"]
        E["epoll 线程<br/>minimio::Poll"]
    end

    OS["OS 事件队列<br/>epoll/kqueue/IOCP"]
    R["Registrator<br/>(主线程持有)"]

    M -->|"Task 经每线程独立通道"| P0
    M -->|"Task"| P1
    P0 -->|"PollEvent::Threadpool(id,cb_id,Js)"| CH["event_reciever<br/>(共享 mpsc 通道)"]
    P1 --> CH
    E -->|"PollEvent::Epoll(token)"| CH
    E -->|"PollEvent::Timeout"| CH
    CH --> M
    R -->|"register(fd, token, READABLE)"| OS
    OS -->|"就绪事件+token"| E
    M -.->|"写入共享 epoll_timeout"| E

11.1 运行输出揭示了什么(TICK 视角)

1
2
3
4
5
6
7
8
9
10
main: ===== TICK 1 =====
main: Immediate1 timed out / Immediate2 timed out ← 0ms 计时器第一轮就到期
pool2/pool3: finished File read ← 线程池与主循环并行工作
main: First count: 39 characters. ← 文件结果回调
main: ===== TICK 2 =====
main: SETTIMEOUT ← 1000ms 计时器
...
epoll: epoll event 7 is ready ← http 响应就绪
...
main: FINISHED

关键观察:

  • 用户代码注册的所有事件先同步注册完,循环才开始
  • 回调在主线程跑,任务在线程池/OS 跑——这就是"并发但不(在用户代码层面)并行"
  • 嵌套回调(文件读完再读下一个)自然产生新的 TICK
  • pending_events 每注册 +1、每回调执行 -1,归零即退出

第 12 章:捷径与改进

书中的简化点:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// 1. 单 tick 回调积压上限: 防止回调太多饿死 I/O 资源
fn run_callbacks(&mut self) {
let mut n = 0;
while n < MAX_CALLBACKS_PER_TICK { /* ... */ n += 1; } // 剩余的下轮再跑
}

// 2. 无线程可用时 panic → 应改为排队等待
fn get_available_thread(&mut self) -> Option<usize> {
self.available_threads.pop() // None 时把任务挂到待发队列
}

// 3. 动态 park 时长: 有积压时不等待/短等待, 避免延迟
// 4. Vec 只增不减(高负载后内存居高不下) → 换 LinkedList 等结构
// 5. 未实现: process.nextTick / setImmediate / 一轮 poll 多事件处理
// / Http 读写全异步 / 伪唤醒重注册 ...

全书核心结论

  1. 同步是错觉:从程序员/OS/CPU 三个参照系看,"顺序执行"只在第一个成立
  2. 并发管效率,并行管吞吐:并发不能让单个任务更快,只能让一组任务总耗时更优
  3. 一切 I/O 皆系统调用:用户态(Ring 3)到内核态(Ring 0)没有捷径
  4. 硬件早已为异步准备好:中断(IDT/IRQ) + DMA + 固件,CPU 全程参与事件通知
  5. 策略 3(OS 事件队列) 是最优解但只解决一半——"如何暂停任务"由运行时决定:Node 用回调,Rust 用状态机 Future
  6. 事件循环 = Timers → Callbacks → Poll(阻塞等待三路事件) 的循环pending_events == 0 时自然退出
  7. "别阻塞事件循环"的机理:所有回调都在主线程串行执行,一个慢回调 = 所有任务停摆
  8. 锁与阻塞的交互是死锁高发区:先 drop 锁再去 poll/recv
  9. 就绪型(epoll/kqueue)与完成型(IOCP)的抽象差异是跨平台异步库的核心难题(借出缓冲区的所有权问题甚至惊动了借用检查器)

延伸阅读:Exploring Epoll, Kqueue and IOCP with RustGreen threads explained in 200 lines of RustExploring Rust Futures


本博客所有文章除特别声明外,均采用 CC BY-SA 4.0 协议 ,转载请注明出处!