
Tokio I/O Driver 重构深度解析从mut self到并发就绪等待的演进之路【免费下载链接】tokioA runtime for writing reliable asynchronous applications with Rust. Provides I/O, networking, scheduling, timers, ...项目地址: https://gitcode.com/GitHub_Trending/to/tokio导读本文基于 Tokio 仓库内的设计文档 tokio/docs/reactor-refactor.md完整还原 Tokio 0.3 版本 I/O Driver 内部重写的设计思路与技术细节。文档记录了 Tokio 如何通过重构Registration与ScheduledIo、引入侵入式 waker 链表、设计 driver tick 竞态防护机制最终让TcpStream等 I/O 类型支持以self发起async fn并发操作。读完本文你将理解 Tokio I/O 就绪通知的底层数据布局、clear_readiness与 tick 的配合原理以及这套设计在 tokio/src/runtime/io/scheduled_io.rs 等源码中的最终落地形态。一、重构背景为什么 0.3 必须重写 I/O Driver在 Tokio 0.3即 1.0 beta之前所有 I/O 类型的async函数都要求mut self。这一限制的根源在于任务的 waker 被存储在 I/O 资源自身的内部状态ScheduledIo中而不是存储在async函数返回的 future 里。由于ScheduledIo每个方向读方向、写方向只能保存一个 waker同一时刻每个方向只能有一个任务在等待。这在实践中意味着无法在同一TcpStream上并发发起多个读操作TcpStream无法直接作为AsyncRead/AsyncWrite使用多个任务共享一个 I/O 资源时waker 会互相覆盖导致通知丢失。设计文档为这次重构列出了明确目标见 tokio/docs/reactor-refactor.md目标让 I/O 类型支持以self调用async fn精炼Registration的 API。非目标实现AsyncRead/AsyncWriteforTcpStream或其他引用类型这一边界在后续设计中再次出现。核心思路一句话概括把 waker 从 I/O 资源的内部状态中移出放入操作返回的 future 中。这样每个操作都能持有独立的 waker从而支持同一资源上注册多个 waiter。文档明确提到Notify中已经验证过的 intrusive wake list侵入式唤醒链表策略可以套用到这一场景但 I/O Driver 还有一些独有难题需要处理。该重构在 tokio/CHANGELOG.md 的 0.3.0 发布说明2020 年 10 月 15 日中得到了印证I/O driver internal rewrite同时net: tcp,udp,uds types support operations with self、io: upgrade to mio 0.7。二、重构Registration用self等待任意 interest 集合Registration在 Tokio 0.3 中被移出公共 API对应 issue #2728CHANGELOG 中记录为 PollEventedandRegistrationare removed降级为支撑TcpStream等 I/O 资源的内部实现细节。设计文档给出了重构后的目标 APIstruct Registration { ... } // TODO: naming struct ReadyEvent { tick: u32, ready: mio::Ready, } impl Registration { /// interest 必须是其他方法中所有 interest 集合的**超集**。 /// 这是传给 mio 的 interest 集合。 pub fn newT(io: T, interest: mio::Ready) - io::ResultRegistration where T: mio::Evented; /// 等待 interest 中包含的任何就绪事件返回表示所收到就绪事件的 /// ReadyEvent。 async fn readiness(self, interest: mio::Ready) - io::ResultReadyEvent; /// 清除由指定 ReadyEvent 表示的资源级就绪状态。 async fn clear_readiness(self, ready_event: ReadyEvent); }关键设计点interest 超集约束new()时传入的 interest 必须覆盖后续所有readiness()调用中可能出现的 interest因为它是唯一传给mio的注册 interest。在 0.3 落地版本中这一概念演化为 tokio/src/io/interest.rs 中的Interest类型READABLE 0b0001、WRITABLE 0b0010配合 tokio/src/io/ready.rs 中细分的READABLE/WRITABLE/READ_CLOSED/WRITE_CLOSED就绪位。self支持并发 waiter不同 waiter 可以携带不同的 readiness interest 并发等待这是从每方向一个 waker到每操作一个 waker的关键转折。新注册即新ScheduledIonew()会在 I/O Driver 中创建一个ScheduledIo条目并把资源注册到mio。在最终代码 tokio/src/runtime/io/registration.rs 中Registration持有scheduler::Handle与ArcScheduledIo提供了readiness(interest)、async_io、poll_read_ready/poll_write_ready等内部方法其中async_io实现了文档中的等待就绪 → 尝试操作 → WouldBlock 则清除就绪 → 循环的完整模式。三、边沿触发与EWOULDBLOCKreadiness 循环的伪代码Tokio 使用边沿触发edge-triggered通知操作系统只在就绪状态发生变化时上报一次事件之后即使资源持续可读/可写也不会再次上报。因此 I/O Driver 必须自行跟踪每个资源当前已知的就绪状态以避免在明知系统调用会返回EWOULDBLOCK时仍然发起无谓的系统调用。文档给出了执行一次 TCP 读取的伪代码这是理解整套设计的最小范例async fn read(self, buf: mut [u8]) - io::Resultusize { loop { // 等待就绪 let event self.readiness(interest).await?; match self.mio_socket.read(buf) { Ok(v) return Ok(v), Err(ref e) if e.kind() WouldBlock { self.clear_readiness(event); } Err(e) return Err(e), } } }执行流程拆解readiness(interest).await先检查当前已知就绪状态与interest是否有交集有交集则立即返回无交集则挂起任务直到 I/O Driver 收到新的就绪事件拿到就绪事件后真正调用read(buf)成功则返回失败且错误为WouldBlock时调用clear_readiness(event)清除本次消费掉的就绪位然后进入下一轮循环等待新事件其他错误直接返回。WouldBlock表示资源已不再就绪此时必须清除已知就绪位否则下一轮readiness()会因为旧就绪位仍然存在而空转甚至反复触发无效系统调用。四、重构ScheduledIo侵入式 waker 链表ScheduledIo是每个 I/O 资源在 Driver 侧的状态载体。重构后它改为使用侵入式 waker 链表链表的每个节点都携带自己传入readiness()的 interest 集合。文档给出的结构如下#[derive(Debug)] pub(crate) struct ScheduledIo { /// 资源已知状态与其他必须原子更新的状态打包在一起 readiness: AtomicUsize, /// 跟踪等待该资源的任务 waiters: MutexWaiters, } #[derive(Debug)] struct Waiters { // 侵入式 waiter 链表 list: LinkedListWaiter, /// 供 AsyncRead 实现使用的 waiter reader: OptionWaker, /// 供 AsyncWrite 实现使用的 waiter writer: OptionWaker, } // 此结构体包含在 readiness() 返回的 **future** 中 #[derive(Debug)] struct Waiter { /// 侵入式链表指针 pointers: linked_list::PointersWaiter, /// 等待 I/O 资源的任务的 waker waiter: OptionWaker, /// 正在等待的就绪事件即传给 readiness() 的值 interest: mio::Ready, /// 不应为 Unpin _p: PhantomPinned, }这一结构与 tokio/src/runtime/io/scheduled_io.rs 中的最终实现几乎一一对应Waiters保留list侵入式LinkedListWaiter、reader、writer三个字段Waiter增加了is_ready标志并保留了PhantomPinned以保证节点在链表中的地址稳定侵入式链表要求节点不被移动。事件分发逻辑当mio上报事件时对应资源的就绪状态被更新随后遍历 waiter 链表——所有interest与收到就绪事件有交集的 waiter 都被唤醒interest 不匹配的 waiter 继续留在链表中等待。在最终代码中这一逻辑由ScheduledIo::wake(ready)实现tokio/src/runtime/io/scheduled_io.rs它使用drain_filter精确摘取满足ready.satisfies(w.interest)的节点并借助WakeList在锁外批量唤醒避免持锁唤醒导致死锁。另一个值得注意的细节ScheduledIo在最终代码中带有repr(align(128))等按架构区分的缓存行对齐x86_64/aarch64/powerpc64 为 128 字节用于避免多核场景下的伪共享false sharing这是文档未展开、但落地时补上的工程优化。五、竞态条件旧事件清除新就绪的死锁文档专门用一节讨论了一个隐蔽的竞态。回到上面的 TCP read 伪代码async fn read(self, buf: mut [u8]) - io::Resultusize { loop { // 等待就绪 let event self.readiness(interest).await?; match self.mio_socket.read(buf) { Ok(v) return Ok(v), Err(ref e) if e.kind() WouldBlock { self.clear_readiness(event); } Err(e) return Err(e), } } }考虑以下时序多个任务并发执行 I/O 操作配合边沿触发语义可能产生死锁任务 A 收到就绪事件调用mio_socket.read(buf)返回WouldBlock在mio_socket.read(buf)返回之后、clear_readiness(event)执行之前一个新的就绪事件到达例如又有数据可读任务 A 执行clear_readiness()把新到达的就绪状态一并清除下一轮循环中readiness().await永远阻塞——因为新就绪事件已经在上一步被错误清除而边沿触发下 OS 不会再次上报。旧方案的缺陷重构前的 I/O Driver 通过总是先注册任务 waker 再执行操作来规避此问题代价是产生大量不必要的任务通知性能不理想。新方案——driver tickI/O Driver 维护一个自增的 tick 值。每次调用mio::poll()都会使 tick 递增每个就绪事件都关联一个 tickDriver 设置资源就绪状态时会把当前 tick 打包进原子usize。readiness()返回的ReadyEvent携带读取就绪值那一刻的 tickclear_readiness()被调用时带上这个ReadyEvent仅当当前 tick 与ReadyEvent中的 tick 一致时才允许清除就绪。若 tick 不一致说明在读到就绪与清除就绪之间 Driver 又轮询到了新事件此时拒绝清除下一轮readiness()不会阻塞并返回包含新 tick 的ReadyEvent。最终代码将这一机制落实在ScheduledIo::set_readinesstokio/src/runtime/io/scheduled_io.rs中Tick::Clear(t)时若tick ! t则直接返回None放弃本次清除Tick::Set时则tick.wrapping_add(1)推进 tickclear_readiness在清除时还会排除READ_CLOSED/WRITE_CLOSED等终结态位因为关闭状态一旦出现就是最终状态不应被消费掉。ScheduledIo的原子位打包布局ScheduledIo的readiness: AtomicUsize被设计为紧凑的位域。文档提出的布局为| shutdown | generation | driver tick | readiness | |-----------------------------------------------| | 1 bit | 7 bits 8 bits 16 bits |其中shutdown和generation在重构前就已存在。文档还注明generation用于处理资源重用问题ScheduledIo在 Driver 的 slab 中可能被复用旧事件需要靠 generation 区分已过时。落地后的最终代码tokio/src/runtime/io/scheduled_io.rs略有演化注释明确写出| shutdown | driver tick | readiness | |----------------------------------| | 1 bit | 15 bits | 16 bits |即READINESS占低 16 位、TICK占中间 15 位、SHUTDOWN占最高 1 位。generation字段在后续演进中被移除或重构文档中的ReadyEvent.tick类型也从草案的u32收敛为u16与 15 位 tick 对应。这种打包方式让读取就绪值 tick可以一次原子加载完成无需加锁是保证并发正确性的关键。六、Drop 时取消兴趣防止 waker 泄漏readiness()返回的 future 把 waker 以侵入式链表节点形式存放在ScheduledIo中。由于readiness()可以被并发调用链表中可能同时存在大量 waker。一旦某个readiness()future 被提前 drop例如select!放弃了该分支必须把对应节点从链表中摘除否则残留的 waker 会导致ScheduledIo持有的引用无法释放形成内存泄漏被唤醒时还会触发对已回收 future 的无效唤醒。文档强调这是侵入式链表方案的硬性要求。最终代码在Readinessfuture 的Drop实现中落实tokio/src/runtime/io/scheduled_io.rsimpl Drop for Readiness_ { fn drop(mut self) { let mut waiters self.scheduled_io.waiters.lock(); // Safety: waiter 只会存放在 waiters 链表中 unsafe { waiters .list .remove(NonNull::new_unchecked(self.waiter.get())) }; } }此外Registration::drop中还会调用clear_wakers()清空reader/writer槽位以打破可能存在的 waker 与Arcdriver::Inner之间的循环引用源码注释引用了 tokio-rs/tokio#3481。七、AsyncRead/AsyncWrite轮询式 API 的专用槽位AsyncRead和AsyncWrite是poll 风格的 trait与readiness()的 future 风格不同这带来两个约束没有与操作关联的 future因此无法用侵入式链表跟踪 waker没有 future 生命周期因此无法在放弃等待时取消 interest。为此ScheduledIo为读写两个方向各保留一个专用 waker 槽位Waiters.reader和Waiters.writer。AsyncRead/AsyncWrite实现只使用这两个槽位存储 waker且不跟踪具体 interest而是假定关注的事件仅有四类Read ready读就绪Read closed读关闭Write ready写就绪Write closed写关闭文档特别指出read closed 与 write closed 只有在 Mio 0.7 中才可用Mio 0.6 时代这两类事件的处理相当混乱a bit messy。这与 0.3.0 CHANGELOG 中 io: upgrade tomio0.7 的记录互为印证。在 tokio/src/io/ready.rs 中可以看到READ_CLOSED/WRITE_CLOSED已作为独立的就绪位存在。也正因每方向只有一个 waker 槽位只能为资源类型本身实现AsyncRead/AsyncWrite不能为Resource实现——为引用类型实现会允许对同一资源发起并发操作而单槽位必然导致 waker 互相覆盖、任务永久挂起死锁。文档还分析了替代方案改用VecWaker存多个 waker 可以解决并发但会引入内存泄漏问题waker 无法在取消时精确移除因此被否决。八、TcpStream的并发读写by_ref提案与split的演进虽然不能为TcpStream实现AsyncRead/AsyncWrite但文档提出了一个等效方案在TcpStream上新增一个返回轻量引用包装的方法impl TcpStream { /// 命名待定Naming TBD fn by_ref(self) - TcpStreamRef_; } struct TcpStreamRefa { stream: a TcpStream, // Waiter 是侵入式 waiter 链表中的节点 read_waiter: Waiter, write_waiter: Waiter, }这样AsyncRead/AsyncWrite可以实现在TcpStreamRefa上每个TcpStreamRef自带read_waiter与write_waiter两个链表节点天然支持并发当TcpStreamRef被 drop 时其关联的所有 waker 资源一并清理。替代split()的用法文档设想有了by_ref()之后TcpStream::split()将不再需要。同一流可以产生两个独立引用放进select!的不同分支同时读写let rd my_stream.by_ref(); let wr my_stream.by_ref(); select! { // 在各自分支中使用 rd 和 wr }也可以把TcpStream放进Arc在多个任务间共享let arc_stream Arc::new(my_tcp_stream); let n arc_stream.by_ref().read(buf).await?;在最终仓库中的落地形态需要说明的是by_ref只是该设计文档的提案方法命名标注为 Naming TBD在 tokio/src/net/tcp/stream.rs 的最终代码中并未以TcpStreamRef形式出现。实际仓库采取了两条路径达成同样的目标self原生方法TcpStream::read/write等操作直接接受selfCHANGELOG 0.3.0 记录 tcp,udp,uds types support operations withself配合ready(interest)、readable()、writable()等self就绪等待方法在同一流上并发读写成为可能split()/into_split()文档中移除所有 split 函数的设想并未发生tokio/src/net/tcp/stream.rs 至今仍保留split(mut self) - (ReadHalf, WriteHalf)与into_split()可static、可跨任务移动作为并发读写的补充手段。这说明设计文档是蓝图最终实现根据 API 打磨CHANGELOG 称之为 APIs are polished and future-proofed做了取舍但其核心机制——ScheduledIo侵入式链表、tick 防竞态、reader/writer 槽位——被完整保留。九、源码落地验证从设计到实现的对照将文档设计与当前仓库代码逐项对照可以确认这套重构的完整落地相关实现集中在 tokio/src/runtime/io/设计文档要点最终代码位置落地情况Registration内部化、selfAPItokio/src/runtime/io/registration.rspub(crate)提供readiness/async_io/clear_readinessScheduledIo侵入式链表tokio/src/runtime/io/scheduled_io.rsWaiters { list, reader, writer }完全对应ReadyEvent { tick, ready }tokio/src/runtime/io/driver.rsReadyEvent { tick: u16, ready: Ready, is_shutdown }driver tick 防竞态ScheduledIo::set_readinessTick::Clear旧 tick 的清除被拒绝注释与文档一致位打包布局tokio/src/runtime/io/scheduled_io.rsshutdown(1) tick(15) readiness(16)AsyncRead/AsyncWrite专用槽位ScheduledIo::poll_readinessWaiters.reader/writer保留Direction::Read/Write两个槽位mio 事件驱动流程tokio/src/runtime/io/driver.rsDriver::turn中set_readiness(Tick::Set, ...)后wake(ready)仓库中还保留了针对 tick 机制的单元测试stale_event_does_not_clear_readiness_after_u8_wraparound位于 tokio/src/runtime/io/scheduled_io.rs验证了即使 tick 经历多次回绕过期的就绪事件也不会错误清除新就绪状态——这正是文档第五节所述竞态的正确性回归测试。十、总结Tokio 0.3 的 I/O Driver 重构本质上是将每资源单 waker的模型升级为每操作独立 waker的模型具体落为四个关键设计waker 归属迁移从ScheduledIo内部状态迁入readiness()future使self并发等待成为可能侵入式唤醒链表ScheduledIo用LinkedListWaiter承载任意数量的并发 waiter事件到达时按 interest 精确唤醒future drop 时精确摘除节点防泄漏driver tick 防竞态每次mio::poll()递增 tick 并随就绪值原子打包clear_readiness只在 tick 匹配时生效杜绝旧清除吞掉新事件的死锁双通道就绪流AsyncRead/AsyncWrite的 poll API 走 reader/writer 专用槽位async fn路径走侵入式链表两种风格各得其所。这份设计文档虽然是 0.3 时代的过去式蓝图但其中对边沿触发、竞态、waker 生命周期等问题的分析至今仍是理解 Tokio I/O 模型的最佳入门材料而它对应的实现代码则继续驱动着当前版本 tokio 的整套异步 I/O 基础设施。【免费下载链接】tokioA runtime for writing reliable asynchronous applications with Rust. Provides I/O, networking, scheduling, timers, ...项目地址: https://gitcode.com/GitHub_Trending/to/tokio创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考