跳转到内容

服务端运行时

源码文件:26 个 · 核对版本 Etemenanki 596916d · katana v3.0.1
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/wake.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/tests/runtime.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/app/src/serve.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/http/protocol.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/trojan/protocol.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/hysteria/protocol.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/datagrams.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/protocols/tests/support/pipeline.rs
  • Etemenanki/environment/tests/integration/udp.rs
  • katana/src/serve.rs

ProxyServerRuntime 驱动一条代理连接。它是一个手写的 Future(按需也可以是 Stream),持有三样东西:面向客户端的传输层、这条连接打开的每个出站,以及 sans-I/O 协议核心。它通过三个固定大小的缓冲区在三者之间搬运字节,只 poll 真正被唤醒的出站,并且从不 spawn task。除 SOCKS 外,每个服务端协议的流量都在这种运行时里运行,所以本页描述的规则决定了 HTTP、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022、Hysteria 2 和 TUN 连接在负载下的行为。

本页面向修改 concepts/src/runtime.rs、concepts/src/buffer.rs 或 concepts/src/wake.rs 的贡献者,以及需要确切知道自己的协议核心何时被调用、有多少空间可用的协议作者。协议核心一侧的约定(ProxyCoreDecode、Event、Effect、Effects sink)见服务端协议核心。本页讲的是调用它的驱动器。

运行时负责:

  • 把传输层读入一个固定的读缓冲区,并把未解析的字节交给协议核心;
  • 在协议核心请求时(Effect::Open),通过构造时传入的 Connector 拨号一个出站;
  • 严格按照协议核心推入的顺序应用 effect,把转发的区间直接写入出站,不做复制;
  • 把出站读入 scratch 缓冲区,并把每个数据块或数据包交给协议核心;
  • 把协议核心 stage 的内容写回传输层;
  • 代协议核心运行一个 deadline 定时器;
  • 统计它搬运的字节数(Traffic);
  • 决定何时不读取。它的全部背压都来自这个决定。

运行时不负责:

  • 解析或生成协议字节。这是协议核心的工作。
  • 选择目的地、路由或解析域名。这些发生在 Connector 及其返回的 link 内部。
  • 维护自己的定时器。唯一的定时器是协议核心通过 Effect::SetDeadline 设置的那一个。
  • 记录日志。出站失败以事件的形式交给协议核心,致命错误以 RuntimeError 的形式交给调用方。
调用方 协议核心 构造函数 模式
app/src/serve.rs → drive HttpCore、TrojanCore、VlessCore、VMessCore、ShadowsocksCore、Ss2022Core 在 accept 得到的 stream 上调用 new showing_progress
katana src/serve.rs → drive TrojanCore、VlessCore、VMessCore、ShadowsocksCore、Ss2022Core 在 accept 得到的 stream 上调用 new showing_progress
protocols/src/hysteria/server/inbound.rs → classifier Hy2StreamCore,每个代理 stream 一个运行时 在 QuicIo 上调用 new quiet
protocols/src/hysteria/server/inbound.rs,每条 QUIC 连接在认证通过后、启用 UDP 时 Hy2UdpCore,覆盖该连接全部数据报的一个运行时 在 QuicDatagrams 上调用 over_datagrams quiet
protocols/src/tun/inbound.rs → serve_stream PassthroughCore,每个 TCP stream 一个运行时 new showing_progress
protocols/src/tun/inbound.rs,每个客户端源地址 TunUdpCore,覆盖该源地址全部 UDP 流的一个运行时 在 TunUdpLink 上调用 over_datagrams quiet

SOCKS 入站有自己的驱动器(socks.serve),不使用这个运行时。

每个生产环境调用方都用协议核心自己的关联常量来实例化 BUF_SIZE,例如 ProxyServerRuntime::<{ Hy2UdpCore::<()>::BUF_SIZE }, _, _, _>::over_datagrams(...),或者在一个对 const BUF: usize 泛型的 drive 中使用 ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector):

协议核心 BUF_SIZE STAGING_RESERVE MAX_DATAGRAM
PassthroughCore(protocols/src/core/mod.rs) 8 * 1024 0 默认 4096
HttpCore MAX_HEAD = 64 * 1024 256 默认 4096
TrojanCore 16 * 1024 由帧开销计算得出 MAX_LENGTH = 8192
VlessCore 16 * 1024 由帧开销计算得出 8192
VMessCore 32 * 1024 4096 8192
ShadowsocksCore 20 * 1024 由帧开销计算得出 默认 4096
Ss2022Core 32 * 1024 由帧开销计算得出 默认 4096
Hy2StreamCore 8 * 1024 2048 默认 4096
Hy2UdpCore 16 * 1024 4096 MAX_UDP_SIZE = 4096
TunUdpCore 8 * 1024 4096 4096
concepts/src/runtime.rs
pub struct ProxyServerRuntime<const BUF_SIZE: usize, Core, Trans, Conn, Mode = ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Conn: Connector<Core::Target>,
{ /* private fields */ }
pub struct ProxyRunsQuiet;
pub struct ProxyShowsProgress;
参数 约束 含义
BUF_SIZE const usize,必须大于 Core::STAGING_RESERVE 三个缓冲区(传输层读缓冲区、传输层 staging 缓冲区、出站 scratch 缓冲区)各自的大小。它同时也是这条连接能接受的最大协议帧的上限。
Core ProxyCoreDecode 协议状态机。它的关联类型确定了出站 key(Key)、拨号目标(Target)、错误类型(Error)和传输层对端地址(TransportAddr)。
Trans 在驱动它的 impl 上为 Transport<Addr = Core::TransportAddr> StreamTransport<T>(来自 new)或 DatagramTransport<D>(来自 over_datagrams)。调用方从不需要写出它。
Conn Connector<Core::Target>,在驱动它的 impl 上还要求 Conn::Datagram: DatagramLink<Addr = Destination> 把一个目标拨号成 Outbound::Stream 或 Outbound::Datagram。
Mode ProxyRunsQuiet(默认)或 ProxyShowsProgress 选择使用 Future impl(resolve 为总的 Traffic),还是 Stream impl(每个工作单元产出一个 Traffic 增量)。
concepts/src/runtime.rs
impl<const BUF_SIZE: usize, Core, T, Conn>
ProxyServerRuntime<BUF_SIZE, Core, StreamTransport<T>, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
T: AsyncRead + AsyncWrite + Unpin,
Conn: Connector<Core::Target>,
{
pub fn new(transport: T, core: Core, connector: Conn) -> Self;
}
impl<const BUF_SIZE: usize, Core, D, Conn>
ProxyServerRuntime<BUF_SIZE, Core, DatagramTransport<D>, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
D: DatagramLink<Addr = Core::TransportAddr>,
Conn: Connector<Core::Target>,
{
pub fn over_datagrams(link: D, core: Core, connector: Conn) -> Self;
}
impl<const BUF_SIZE: usize, Core, Trans, Conn>
ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Conn: Connector<Core::Target>,
{
pub fn showing_progress(
self,
) -> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>;
}
  • new 把 stream 包装成 StreamTransport,over_datagrams 把 link 包装成 DatagramTransport。两者随后都调用私有的 build,它分配三个缓冲区并设置 prefer_transport = true,因此第一次 poll 会先尝试传输层,再尝试出站。

  • build 一开始就执行:

    assert!(
    BUF_SIZE > Core::STAGING_RESERVE,
    "BUF_SIZE must exceed the core's STAGING_RESERVE or nothing is ever read"
    );

    这是运行期的 assert!,所以错误的实例化会在构造运行时的时候 panic,而不是在编译时报错。原因是:只有 staging 空间达到 STAGING_RESERVE + 1 字节时才会读取 stream 出站,而不大于预留量的缓冲区永远达不到这个条件;另外 datagram_limit 要计算 BUF_SIZE - Core::STAGING_RESERVE。

  • showing_progress 消耗 quiet 形式的运行时,把所有字段移入 ProxyShowsProgress 形式。它只在 quiet 类型上可用,所以一个运行时最多切换一次模式。所有调用方都在构造后立即调用它。

concepts/src/runtime.rs
impl<const BUF_SIZE: usize, Core, Trans, Conn, Mode>
ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, Mode>
where
Core: ProxyCoreDecode,
Conn: Connector<Core::Target>,
{
pub fn core(&self) -> &Core;
pub fn summary(&self) -> Traffic;
}
// On the impl that drives the runtime
// (Trans: Transport<Addr = Core::TransportAddr>, Conn::Datagram: DatagramLink<Addr = Destination>):
pub const fn datagram_limit() -> usize;
  • core() 让调用方在两步之间检查协议核心。服务层用它询问握手是否已经结束(is_established())。
  • summary() 返回到目前为止搬运的字节数。
  • datagram_limit() 是数据报出站能完整交付的最大数据包:取 Core::MAX_DATAGRAM,上限为 BUF_SIZE - Core::STAGING_RESERVE。对 TunUdpCore 来说是 min(4096, 8192 - 4096) = 4096 字节。
  • 只要 Trans: Unpin,该类型就实现 Unpin。传输层是 Unpin 的,连接 future 和定时器都经过 box,所以内部没有任何东西依赖被 pin 住。

客户端一侧通过同一个 trait 来看待,因此运行时只需为两种传输层写一次。这个 trait 是公开的,但调用方从不需要写出它。

concepts/src/runtime.rs
pub enum Received<A> {
Bytes,
Datagram(A),
Eof,
}
pub trait Transport: Unpin {
const DATAGRAM: bool;
type Addr;
fn poll_recv(
&mut self,
cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<io::Result<Received<Self::Addr>>>;
fn poll_send(
&mut self,
cx: &mut Context<'_>,
data: &[u8],
to: Option<&Self::Addr>,
) -> Poll<io::Result<usize>>;
fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
fn poll_shutdown(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
}
pub struct StreamTransport<T>(pub T);
pub struct DatagramTransport<D>(pub D);
StreamTransport<T> DatagramTransport<D>
约束 T: AsyncRead + AsyncWrite + Unpin D: DatagramLink
DATAGRAM false true
Addr () D::Addr:对于数据包带有对端地址的 link(服务多个对端的 UDP socket,或 TUN 按源地址划分、以远端地址寻址的 TunUdpLink)是 SocketAddr;对于单对端 link(例如 QUIC 连接的数据报)是 ()
poll_recv 什么都没读到的读取是 Received::Eof,否则是 Received::Bytes 总是 Received::Datagram(from)。没有流结束。
poll_send poll_write;忽略 to poll_send_to(to)。to = None 是 InvalidInput 错误,“a datagram transport needs a peer per packet”。
poll_flush、poll_shutdown 转发给 T 空操作,返回 Ready(Ok(()))

出站一侧使用 concepts/src/link.rs 中的 trait:

concepts/src/link.rs
pub trait Connector<Target> {
type Stream: AsyncRead + AsyncWrite + Unpin;
type Datagram: DatagramLink;
type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;
}
pub enum Outbound<S, D> {
Stream(S),
Datagram(D),
}

任何 future 产出 Outbound 的 FnMut(Target) -> Fut 都是 Connector,测试就是这样传入闭包的。DatagramLink、UdpOutbound 和 Destination 见 Link 与类型。

连接在 build 中一次性分配三个 BUF_SIZE 字节的 boxed 数组,之后从不扩容。buffer::boxed_array 直接在堆上构造每个数组,所以即使 BUF_SIZE 很大也不会经过栈。被代理的字节所经历的每一次复制都落在这些缓冲区中的某一个里,而背压就来自它们被填满。

flowchart LR
  client["客户端传输层"]
  up["up: ReadBuffer"]
  core["Core::handle"]
  queue["effect 队列"]
  out["出站 link"]
  scratch["scratch"]
  staging["staging: WriteBuffer"]
  client -->|"poll_recv"| up
  up -->|"Event::Transport"| core
  core -->|"Forward 区间"| queue
  queue -->|"poll_write, poll_send_to"| out
  out -->|"poll_read, poll_recv_from"| scratch
  scratch -->|"Event::Outbound, Event::Datagram"| core
  core -->|"stage, reserve, stage_to"| staging
  staging -->|"poll_send"| client

上行字节完全不经运行时复制。协议核心在读缓冲区中原地解密,并转发其中的区间,运行时直接从 up 把它们写到出站。下行字节被读入 scratch,在协议核心把它们封装进 staging 时再复制一次。

传输层被读入 up,它分为三个区域:

0 ............... start ............... end .............. BUF_SIZE
| 已解析 | 未解析 | 空闲 |
| (排队中的 | (被拆开的帧 | (下一次读取 |
| 转发可能仍 | 的尾部) | 落在这里) |
| 在读取它) | | |
方法 签名 用途
start pub fn start(&self) -> usize 把传输层事件中的转发区间转换为绝对偏移时使用的基址
unparsed pub fn unparsed(&mut self) -> &mut [u8] 作为 Event::Transport 交给协议核心的切片;可变,以便协议核心原地解密
unparsed_len pub fn unparsed_len(&self) -> usize 判断是否还有东西需要交给协议核心
slice pub fn slice(&self, range: Range<usize>) -> &[u8] 应用排队中的转发时,解析它的绝对区间
free_len、free pub fn free_len(&self) -> usize、pub fn free(&mut self) -> &mut [u8] 下一次读取写入的尾部
advance_end pub fn advance_end(&mut self, n: usize) 记录一次读到的 n 字节
advance_start pub fn advance_start(&mut self, n: usize) 记录协议核心消费的字节数
compact pub fn compact(&mut self) 把未解析的尾部移到偏移 0。start == 0 时是空操作,没有未解析数据时直接重置。
is_saturated pub fn is_saturated(&self) -> bool start == 0 && end == N:未解析区域占满了整个数组

排队中的转发持有指向这个数组的绝对偏移,所以只要还有转发在排队,已解析区域就不能移动。因此压缩是惰性的,并且有保护条件:

  • Stream 传输层。 poll_transport_read 只在空闲尾部为空(free_len() == 0)时压缩。如果此时缓冲区已饱和,说明协议核心还需要一帧中放不下的更多数据,连接以 RuntimeError::FrameTooLarge 失败。只要有来自传输层的转发在排队,读取路径就会提前返回,所以压缩永远不会移动转发仍在引用的字节(见来源锁定)。
  • 数据报传输层。 读取前缓冲区总是空的,因为一个数据包会被整体消费,而仍在读取它的转发会阻止下一次读取。因此 poll_transport_recv_packet 在每次读取前都会压缩(即重置),并用 debug_assert_eq! 断言没有未解析数据。

发往传输层的字节先 stage 在 staging 中。协议核心通过 Staging 写入器在 end 处追加,运行时把 [start .. end) 写出,并向前推进 start。

0 ............. start ............. end ............. BUF_SIZE
| 已写出 | 待写出 | 剩余空间 |
concepts/src/buffer.rs
impl<const N: usize> WriteBuffer<N> {
pub fn is_empty(&self) -> bool;
pub fn pending(&self) -> &[u8];
pub fn advance_start(&mut self, n: usize);
pub fn room(&mut self) -> usize;
pub fn staging(&mut self) -> Staging<'_>;
}
impl Staging<'_> {
pub fn room(&self) -> usize;
pub fn reserve(&mut self, len: usize) -> Option<&mut [u8]>;
pub fn put(&mut self, bytes: &[u8]) -> Option<()>;
}
  • 所有待写出的字节都写完后,advance_start 把两个偏移都重置为 0,所以排空的缓冲区总能提供完整的大小。
  • 当把待写出字节滑到最前面才能腾出空间时(即 start != 0 且尾部已满,end == N),room() 会这样做。这是该缓冲区自行进行的唯一一次复制。正因如此 room() 接受 &mut self,而 staging() 会先调用它,保证协议核心总能拿到最大的可用尾部。
  • Staging::reserve 在尾部占用 len 字节并返回它们,让协议核心原地填充;否则返回 None,不占用任何空间。put 就是 reserve 加一次复制。帧在原处封装,不经过中间的 Vec。

scratch: Box<[u8; BUF_SIZE]> 接收出站读到的数据。service_key 最多读取:

  • 对 stream 出站,(staging.room() - Core::STAGING_RESERVE).min(BUF_SIZE) 字节;
  • 对数据报出站,datagram_limit() 字节。

第一个上限保证了事件约定成立:一个 n 字节的出站数据块只会在 staging 至少有 STAGING_RESERVE + n 字节空闲时交付,所以协议核心能把它全部封装。

字段 类型 用途
outbounds BTreeMap<Core::Key, Slot<...>> 每个存活的出站 key 一个槽位
ready ReadyQueue<Core::Key> 其出站已发出就绪信号的 key
drained VecDeque<Core::Key> 已从 ready 取出、尚未处理的 key,最早的在前
starved Vec<Core::Key> 需要处理、但因 staging 空间不足或缓冲区被锁定而被搁置的 key
packets PacketList<Core::TransportAddr>,即 VecDeque<(usize, A)> 在数据报传输层下,staging 中的数据包,队首在前
effects EffectList<Core>,即 SmallVec<[Effect<C>; INLINE_EFFECTS]> 协议核心推入 effect 的 sink;两次调用之间为空
queue SmallVec<[Queued<Core>; INLINE_EFFECTS]> 尚未应用的 effect,队首在前
deadline、deadline_armed Option<Pin<Box<Sleep>>>、bool 唯一的定时器
delta、summary Traffic 当前工作单元中搬运的字节数,以及总字节数

标志位 up_needs_feed、transport_read_closed、transport_write_closed、transport_needs_flush、shutdown_transport、finished、failed 和 prefer_transport 构成其余的状态。除三个缓冲区外,每打开一个出站要花费一个 Arc<KeyWaker> 和一个 boxed 的拨号 future;定时器在第一次 SetDeadline(Some(_)) 时 box 一次。

每个存活的 key 在 outbounds 中拥有一个 Slot:

concepts/src/runtime.rs
enum LinkState<F, S, D> {
Connecting(Pin<Box<F>>),
Stream(S),
Datagram(D),
}
struct Slot<K, F, S, D> {
link: LinkState<F, S, D>,
waker: Arc<KeyWaker<K>>,
read_closed: bool,
write_closed: bool,
need_flush: bool,
}

当一次转发已经写入、但 link 的 poll_flush 返回 Pending 时,设置 need_flush。service_key 在再次读取该 key 之前会先完成这次 flush。

stateDiagram-v2
  [*] --> Connecting: Open effect 到达队首
  Connecting --> Stream: 拨号返回 Outbound Stream
  Connecting --> Datagram: 拨号返回 Outbound Datagram
  Connecting --> [*]: 拨号失败,ConnectFailed
  Connecting --> [*]: Close effect
  Stream --> [*]: Close effect
  Stream --> [*]: 写入、flush、shutdown 或读取出错,OutboundError
  Stream --> [*]: 已 shutdown 且读到 EOF
  Datagram --> [*]: Close effect
  Datagram --> [*]: 接收出错,OutboundError
  • 打开。 如果该 key 已经有槽位(包括仍在连接中的槽位),Effect::Open 以 DuplicateKey 使连接失败。否则运行时用 ready.waker(key) 生成一个 waker,调用 connector.connect(target),把 future box 起来存为 Connecting,并调用 schedule(),让拨号被 poll 一次。拨号在 Open 到达队首时才开始,而不是在协议核心推入它的时候:排在一个停滞的转发后面的 Open 需要等待。
  • 连接。 拨号 resolve 后,槽位变为 Stream 或 Datagram,该 key 被再次调度,协议核心收到 Event::Connected { key }。拨号失败会移除该 key(见应用 effect 中的 forget_key),并交付 Event::ConnectFailed { key, error }。
  • 半关闭。 stream 槽位在两个方向都关闭后被移除。Effect::Shutdown 在 poll_shutdown 完成后设置 write_closed;读到零字节会设置 read_closed 并交付 Event::OutboundEof。两者中后发生的那个会移除槽位。
  • 数据报槽位。 对数据报槽位的 Effect::Shutdown 立即完成,只标记写方向。数据报 link 没有流结束,所以槽位只能通过 Close 或接收错误结束。
  • 关闭。 Effect::Close 在任何状态下都立即移除槽位,并 drop 掉 link 或进行中的拨号 future。它不会清除排在它后面、针对同一 key 的 effect,所以之后发往该 key 的转发会以 UnknownKey 失败,除非在那之前有针对该 key 的 Open。

一个 task 要 poll 许多出站。如果每次唤醒都 poll 全部出站,每次唤醒的代价是 O(N),因此每个出站都用自己的 waker 来 poll,由它记录哪个 key 就绪了。这与 FuturesUnordered 的做法相同,只是没有它的侵入式链表。

concepts/src/wake.rs
struct Shared<K> {
ready: Mutex<Vec<K>>,
parent: Mutex<Option<Waker>>,
}
pub struct ReadyQueue<K> {
shared: Arc<Shared<K>>,
}
impl<K> ReadyQueue<K> {
pub fn new() -> Self;
pub fn register(&self, waker: &Waker);
pub fn waker(&self, key: K) -> Arc<KeyWaker<K>>;
pub fn drain_into(&self, into: &mut VecDeque<K>);
pub fn is_empty(&self) -> bool;
}
pub struct KeyWaker<K> {
key: K,
queued: AtomicBool,
shared: Arc<Shared<K>>,
}
impl<K: Copy> KeyWaker<K> {
pub fn key(&self) -> K;
pub fn clear_queued(&self);
}
impl<K: Copy + Send + Sync + 'static> KeyWaker<K> {
pub fn waker(self: &Arc<Self>) -> Waker;
pub fn schedule(self: &Arc<Self>);
}
impl<K: Copy + Send + Sync + 'static> Wake for KeyWaker<K> { /* wake, wake_by_ref */ }

两个 mutex 都是 parking_lot::Mutex。由于 KeyWaker 实现了 std::task::Wake,把它变成 Waker 只是一次引用计数加一,不是一次分配;唯一的分配是每次打开时生成的那个 Arc。

sequenceDiagram
  participant IO as 出站 I/O 驱动
  participant KW as 该 key 的 KeyWaker
  participant RQ as ReadyQueue
  participant T as 运行时 task
  T->>RQ: 每次 poll_once 注册 task waker
  IO->>KW: wake_by_ref
  KW->>KW: queued.swap(true)
  KW->>RQ: 推入 key,仅当它尚未入队
  KW->>RQ: 取走 parent waker
  KW->>T: 唤醒 parent
  T->>RQ: drain_into(drained),最早的在前
  T->>KW: clear_queued
  T->>IO: 用该 key 自己的 waker poll 出站
规则 机制
一个 key 无论被唤醒多少次,最多入队一次 只有当 queued.swap(true, Ordering::AcqRel) 返回 false 时,wake_by_ref 才推入该 key
一连串唤醒只唤醒 task 一次 wake_by_ref 会 take() 走 parent waker。poll_once 每次 poll 都调用 register 把它放回;如果已存的 waker 与新的 will_wake,register 会跳过 clone。
poll 期间到达的唤醒不会丢失 service_key 在 poll 出站之前调用 clear_queued()(一次 Release store),所以 poll 期间的唤醒会让该 key 再次入队,而不是被去重吞掉
新 key 会被 poll 一次 waker(key) 返回的 waker 处于未入队状态,所以运行时在 Open 之后调用 schedule()
需要在没有 I/O 事件时 poll 的 key 能得到 poll schedule() 就是 wake_by_ref()。运行时在 Open 之后、拨号完成之后、每次读到数据或数据包之后(就绪 waker 只在返回 Pending 之后才会触发),以及释放一个饥饿的 key 时调用它。由于它同时会唤醒 parent,poll 期间被调度的 key 一定会再得到一次 poll。

poll_outbounds 只在 drained 为空时才从 ready 补充它。随后它从前往后处理各个 key,并在第一个取得进展的 key 之后返回。读取成功的 key 会被重新调度到就绪队列末尾,所以繁忙的出站轮流得到处理。返回 Pending 的 key 会从 drained 中丢弃;等它的 I/O 就绪时,它自己的 waker 会让它再次入队。已经没有槽位的 key(Close 之后的过期唤醒)会被跳过。

协议核心在 handle 期间把 effect 推入 effects sink。每次调用之后,运行时校验这些 effect 并把它们移到 queue 上,同时给每个 effect 标记其区间所索引的缓冲区:

concepts/src/runtime.rs
enum Source {
Transport,
Scratch,
Held,
}
struct Queued<C: ProxyCoreDecode> {
effect: Effect<C>,
source: Source,
}

enqueue(source, base, limit) 检查每个区间并把它转换为绝对区间:

  • Forward 和 SendTo 的区间必须满足 start <= end <= limit,其中 limit 是该事件字节切片的长度。随后它们按 base 平移:传输层事件为 up.start(),出站事件为 0。不携带字节的事件以 limit = 0 入队,所以只能转发空区间。
  • ForwardHeld 和 SendToHeld 的区间必须满足 start <= end <= core.held().len()。它们保持相对形式,在应用时再次针对 held() 解析,所以两次之间 Vec 重新分配也没有问题。
  • 任何违规都会以 RangeOutOfBounds 使连接失败。

排队中的带区间 effect 会锁定它所读取的缓冲区:在该 effect 被应用之前,不会有新数据读入那个缓冲区。Queued::pins(source) 针对单个 effect 回答这个问题,pinned(source) 针对整个队列回答。

排队中的 effect 锁定 排队期间
来自 Event::Transport 或 Event::TransportDatagram 的 Forward 或 SendTo Source::Transport 不读取传输层,也不再向协议核心交付更多传输层字节(改为设置 up_needs_feed)
来自 Event::Outbound 或 Event::Datagram 的 Forward 或 SendTo Source::Scratch 不处理任何已连接的出站;被唤醒的 key 暂存到 starved
ForwardHeld 或 SendToHeld Source::Held 不读取也不交付传输层数据,不处理任何已连接的出站
Open、Shutdown、Close、ShutdownTransport、SetDeadline、Finish 无 没有限制

held 锁定让协议核心可以安全地改写自己的 held 缓冲区。只要有任何 held effect 在排队(held_free() 为 false),协议核心就收不到任何字节事件:没有 Transport、TransportDatagram、Outbound、Datagram、OutboundEof 或 TransportEof。仍处于 Connecting 状态的槽位不受 scratch 和 held 锁定影响,因为 poll 拨号不会触及这两个缓冲区,所以 Connected 和 ConnectFailed 仍会到达;Deadline、TransportSendFailed 以及应用 effect 期间产生的失败事件也一样。这个锁定只有在协议核心的调用返回、它的 effect 入队之后才开始。在同一次调用内,先前推入的 held 区间在调用结束后仍然读取 held(),所以协议核心在下一个字节事件之前必须保持这些字节原地不动。protocols/src/mux/demux.rs → Demux::feed_chunks 遵循这条规则:VMessCore 在一次调用中把一次读取所解开的全部 chunk 都交给它,而 demux 只在处理下一次读取的开头才裁剪自己的 held 缓冲区。面向协议核心作者的规则见服务端协议核心。

drive_effects 从前往后应用队列,在第一个无法完成的 effect 处停下。它返回是否应用了任何 effect。出站写入用该 key 自己的 waker 来 poll,而不是用 task 的 waker,所以一个被阻塞的转发只会唤醒它自己的 key。

Effect 应用方式 阻塞条件 错误与事件
Forward、ForwardHeld 空区间直接丢弃。否则对该区间 poll_write,然后 poll_flush。部分写入会推进 range.start,effect 留在队首。 槽位处于 Connecting;poll_write 返回 Pending。flush 挂起会设置 need_flush,但不阻塞。 没有槽位:UnknownKey。数据报槽位:WrongLinkKind。写入 0 字节(WriteZero)、写入错误或 flush 错误:丢弃该 key,协议核心收到 OutboundError。
SendTo、SendToHeld 一次 poll_send_to(data, &to) 槽位处于 Connecting;Pending 没有槽位:UnknownKey。stream 槽位:WrongLinkKind。发送被拒绝:丢弃该数据包,key 保持存活,协议核心收到 SendFailed。
Open 见出站槽位 从不 key 仍存活:DuplicateKey
Shutdown 对 stream 执行 poll_shutdown;对数据报槽位立即完成 槽位处于 Connecting;Pending 没有槽位:UnknownKey。I/O 错误:丢弃该 key,协议核心收到 OutboundError。
Close 移除槽位 从不 无,即使 key 不存在
ShutdownTransport 设置 shutdown_transport 从不 无
SetDeadline(after) Some:把 boxed 的 Sleep 重置为 tokio::time::Instant::now() + after(首次时创建)并启用。None:停用。 从不 无
Finish 设置 finished 从不 无

当一个 key 丢失时,forget_key 移除其槽位,以及所有针对它排队的 Forward、SendTo、ForwardHeld、SendToHeld、Shutdown 和 Close。等待一个刚刚失败的拨号的转发已无处可去,而协议核心会通过恰好一个事件(ConnectFailed 或 OutboundError)得知这次丢失。协议核心的应对被入队到队列末尾,排在仍在等待的所有内容之后。

在一轮应用了 effect 的处理结束时,如果有饥饿的 key,而且 scratch 和 held 缓冲区都不再被锁定,就重新调度这些饥饿的 key。转发和发送写出的字节计入 delta.outbound_tx。

poll_once 把 task 的 waker 注册到 ReadyQueue,然后调用 step,后者最多执行一个工作单元。step 返回 Some(true) 表示有进展,Some(false) 表示在被唤醒之前无事可做,连接结束后返回 None。

flowchart TB
  poll["poll_once: ready.register(task waker)"]
  dl{"已启用的 deadline 触发了?"}
  fx{"drive_effects 应用了 effect?"}
  wr{"poll_transport_write 取得进展?"}
  done{"已 finished、staging 为空、shutdown 完成?"}
  fin{"已 finished?"}
  io["读传输层和处理出站,优先的一侧先来"]
  moved{"有一侧取得进展?"}
  prog["有进展:把 delta 移入 summary"]
  complete["完成"]
  pend["Pending"]
  poll --> dl
  dl -->|是| prog
  dl -->|否| fx
  fx -->|是| prog
  fx -->|否| wr
  wr -->|是| prog
  wr -->|否| done
  done -->|是| complete
  done -->|否| fin
  fin -->|是| pend
  fin -->|否| io
  io --> moved
  moved -->|"是,下次优先未取得进展的一侧"| prog
  moved -->|否| pend

依次为:

  1. Deadline。 poll_deadline poll 已启用的 Sleep。定时器排在最前,这样 I/O 总是就绪的连接也不会饿死自己的握手或空闲 deadline。触发时,定时器被停用,协议核心收到 Event::Deadline。
  2. Effect。 drive_effects 重试仍在排队的内容,通常是一个在等待出站或拨号的转发。
  3. 传输层写入。 poll_transport_write 写出 stage 的字节:stream 上每个 step 一次 poll_send,数据报传输层上则发送所有可发送的数据包。staging 变空后它执行 flush,然后执行被请求的 poll_shutdown。
  4. 完成检查。 当 finished 已设置、staging 为空,并且(如果请求了 ShutdownTransport)传输层的写方向已关闭时,运行时完成。
  5. 交替读取。 除非协议核心已经 finished,step 按 prefer_transport 给出的顺序尝试 poll_transport_read 和 poll_outbounds。传输层读取取得进展会设置 prefer_transport = false;出站取得进展会把它设回 true。任何一侧都不会饿死另一侧。
  6. 无事可做。 该 step 报告没有进展,poll_once 返回 Pending。此时每个被 poll 过的来源都已注册了 waker:定时器和传输层用 task 的 Context,每个出站用自己的 key waker,后者通过 ReadyQueue 回到 task。

取得进展的工作单元把 delta 移入 summary,并把 delta 交给对应模式的 impl。

WORK_BUDGET。 quiet 形式的 Future::poll 最多循环执行 WORK_BUDGET(64)个工作单元。如果预算用完时每个工作单元仍在取得进展,它会调用 cx.waker().wake_by_ref() 并返回 Pending。因此一条传输层或出站总是同步就绪的繁忙连接无法独占 executor 的 worker。stream 模式每次 poll_next 只执行一个工作单元,所以由调用方控制节奏。

stream 传输层上的 poll_transport_read:

  1. 如果读方向已关闭、Source::Transport 或 Source::Held 被锁定,或者 staging.room() < Core::STAGING_RESERVE,则不取得进展直接返回。
  2. 如果设置了 up_needs_feed,先把剩下的未解析字节交给协议核心;如果这消费了数据或再次锁定了缓冲区,就报告有进展。
  3. 如果空闲尾部为空:缓冲区已饱和时以 FrameTooLarge 失败,否则压缩。
  4. 读入空闲尾部。流结束时设置 transport_read_closed 并交付 Event::TransportEof。读到的字节计入 transport_rx 并交给协议核心。

feed_transport 在循环中交付 Event::Transport(up.unparsed())。每次调用后,它检查消费的字节数,以 base = up.start() 把 effect 入队,推进 start,然后应用这些 effect。当没有未解析数据、协议核心什么都没消费,或者出现了锁定(此时设置 up_needs_feed)时停止。所以只要协议核心还在消费,就会再次调用它;只有当协议核心消费完所有数据或停止消费后,运行时才会读取更多数据。staging 空间只在读取之前检查一次,而不是在这个循环的每次调用之前检查。

service_key(key) 对一个出站 poll 一次:

  1. 如果该 key 没有槽位,直接返回。否则调用 clear_queued(),并用该 key 的 waker 构造一个 Context。
  2. 如果设置了 need_flush,先完成 flush。Pending 则不取得进展直接返回;出错则使该 key 失败。
  3. 计算该 key 所需的 staging 空间:数据报槽位为 STAGING_RESERVE + datagram_limit(),其他情况(包括仍在连接中的槽位)为 STAGING_RESERVE + 1。如果空间不足,或者 scratch 或 held 缓冲区被锁定且槽位不处于 Connecting,就把该 key 推入 starved 并返回。
  4. Connecting:poll 拨号(见出站槽位)。
  5. Stream:除非 read_closed,读入 scratch[..max]。读到零字节表示流结束。否则重新调度该 key,计入 outbound_rx,协议核心收到 Event::Outbound { key, data },并且必须整体消费。
  6. Datagram:接收一个数据包到 scratch[..datagram_limit()]。重新调度该 key,协议核心收到 Event::Datagram { key, from, data },并且必须整体消费。接收错误使该 key 失败;这是数据报出站不经 Close 而丢失的唯一方式。

缓冲区从不扩容,所以每条背压规则的形式都相同:在结果有地方可去之前,运行时拒绝读取。没有读取的字节留在内核 socket 或对端的发送窗口中,其余的由对端自己的流量控制完成。

条件 原因
transport_read_closed 传输层已经返回了流结束
有来自传输层的转发在排队(pinned(Source::Transport)) 它的字节仍以绝对偏移留在读缓冲区中;读取和压缩可能移动它们。因此停滞的出站会让上行停滞,而不是把上行数据缓冲起来。
有 held effect 在排队(!held_free()) held 缓冲区被锁定期间,协议核心不能收到字节事件
staging.room() < Core::STAGING_RESERVE 协议核心必须能够 stage 它的应答。停止读取的客户端会填满 staging,然后它自己的上行也会停止。
finished 协议核心 finished 之后,step 不再读取

当运行时因为锁定而停止交付时,它会设置 up_needs_feed,下一次读取尝试会先交付剩下的未解析字节,再读取新数据。

出站 只在 staging 空间达到以下值时处理 读取大小
Stream(或仍在连接中) STAGING_RESERVE + 1 room - STAGING_RESERVE,上限为 BUF_SIZE
数据报 STAGING_RESERVE + datagram_limit() datagram_limit()

数据报出站会等到有足够容纳一个完整数据包的空间,绝不接受更少。因此停止读取的客户端会让下行停滞在出站 socket 处,而不是收到被截断的数据包。

在运行时无法接收其数据时被唤醒的 key(见服务一个出站 key 的第 3 步)会被推入 starved,不被 poll。它的入队标志已经清除,所以它自己的 waker 仍能让它入队。reschedule_starved 清空 starved,并对每个仍有槽位的 key 调用 schedule()。它在以下时机运行:

  • stream 传输层的一次写入至少接受了一个字节,从而腾出了 staging 空间之后;
  • 数据报传输层的发送循环发送或丢弃了至少一个数据包之后;
  • 一轮应用了 effect 的 drive_effects 结束时,只要 scratch 和 held 缓冲区都不再被锁定。

被重新调度但空间仍然不足的 key 只会再次进入饥饿状态。

一个停滞的写入会卡住整个队列

Section titled “一个停滞的写入会卡住整个队列”

effect 严格按顺序应用,一个无法完成的转发会卡住它后面的所有 effect。对只有一个出站的协议核心来说,这正是想要的背压。对多路复用的协议核心来说,这意味着一个缓慢的子流会卡住所有兄弟子流的上行,直到它排空为止,因为它来自传输层的转发锁定了读缓冲区,传输层不会被读取。来自其他出站的下行照常流动:它们读到的数据落在 scratch 中,并直接封装进 staging,不需要等待队列。

sequenceDiagram
  participant C as 客户端传输层
  participant R as 运行时
  participant A as 出站 1,缓慢
  participant B as 出站 2
  C->>R: key 1 的 DATA,100 字节
  R->>A: poll_write 接受 8 字节,然后 Pending
  Note over R: Forward 留在队首,读缓冲区被锁定
  C--)R: FIN,转发排队期间不读取
  B->>R: 下行字节
  R->>C: 经 staging 发送封装好的 key 2 的 DATA
  A-->>R: key waker 触发,再次可写
  R->>A: 剩余 92 字节
  R->>C: 恢复读取,解析 FIN

concepts/tests/runtime.rs 中的 stalled_outbound_holds_uplink_but_not_other_downlink 固定了这个场景。

运行时在读取之前检查 staging 空间,以保证协议核心能封装交给它的内容:

引出该事件的 poll 事先检查的空间
传输层读取:Transport、TransportDatagram、TransportEof STAGING_RESERVE
stream 出站处理:Outbound、OutboundEof,以及拨号的 Connected 或 ConnectFailed STAGING_RESERVE + 1。n 字节的数据块只会被读入 room - STAGING_RESERVE 的空间,所以它到达时有 STAGING_RESERVE + n 空闲。
数据报出站处理:Datagram STAGING_RESERVE + datagram_limit()
Deadline 无
应用 effect 或发送数据包期间产生的失败:写入、flush 或 shutdown 产生的 OutboundError、SendFailed、TransportSendFailed 无

无论空间如何,Deadline 都会交付,这样即使对端停止读取,空闲超时仍然可以触发。最后一行的失败事件由遇到失败的那一轮处理产生,不会单独检查空间。协议核心在处理这些事件时必须容忍 stage 或 reserve 返回 None。更一般地说,协议核心应当把 staging 写入器返回的 None 转成它返回的错误(测试用的协议核心写成 .ok_or("staging full")?),绝不能把它当作不可达。

传输层字节预期被转发,而不是被 stage。预留量只在读取传输层之前检查一次,之后 feed_transport 可能会针对同一批字节多次调用协议核心,所以只有第一次调用能确定预留空间是空闲的。回显传输层字节的协议核心应读取 Staging::room,在应答放不下时少消费一些字节。

运行时在每次调用后检查消费的字节数:

事件 协议核心必须消费
Transport(data) 至多 data.len()。剩余部分留在读缓冲区中,连同新读到的字节再次提供。
TransportDatagram、Outbound、Datagram 恰好 data.len()。剩余部分无处保存。
不携带字节的事件 任意;返回值被忽略

违反约定会以 BadConsume 使连接失败。

over_datagrams 在一个 DatagramLink 上构建运行时:可以是数据包带有对端地址的 link(TransportAddr = SocketAddr),例如服务多个对端的一个 UDP socket 或 TUN 的 TunUdpLink;也可以是 QUIC 连接的数据报一侧(TransportAddr = ())。协议核心通常充当解复用器。例如 Hy2UdpCore 按会话 id 为出站设置 key,第一次见到某个 id 时打开该会话的出站,并通过 deadline 关闭空闲会话。TunUdpCore 对整个源地址只使用一个出站 key,并把每个应答 stage 到它来源的远端地址,TunUdpLink 再把该地址映射回客户端的流。

  • 读取遵循内核 socket 的规则:每次读取一个数据包,读入空的读缓冲区,最多有 Core::MAX_DATAGRAM.min(BUF_SIZE) 字节的空间。更长的数据包会被 link 的接收截断。数据包以 Event::TransportDatagram { from, data } 交付,必须整体消费。如果 link 返回了 Received::Bytes 或 Received::Eof,连接以 RuntimeError::Transport(InvalidData,“datagram transport delivered a stream read”)失败。

  • Staging 要求每个字节都有对端。deliver 用 Effects::with_packets 构造 sink,协议核心使用以下方法 stage:

    concepts/src/core.rs
    pub fn stage_to(&mut self, to: C::TransportAddr, len: usize) -> Option<&mut [u8]>;
    pub fn put_to(&mut self, to: C::TransportAddr, bytes: &[u8]) -> Option<()>;
    pub fn transport_is_datagram(&self) -> bool;

    每次 stage_to 预留 len 字节,并把 (len, to) 推入 packets。协议核心返回后,deliver 比较 staging 缓冲区的增长量与新数据包长度之和。用普通 stage stage 的字节会破坏这个等式,连接以 StagedWithoutPeer 失败。在 stream 传输层下情况正好相反:stage_to 返回 None,transport_is_datagram() 为 false。

  • 发送在 poll_transport_send_packets 中逐个数据包进行:对队首数据包执行 poll_send(&pending[..len], Some(&to)),然后弹出该数据包,staging 推进 len。Pending 会停止循环。被 link 拒绝的数据包会被丢弃,不计入 transport_tx,并以 Event::TransportSendFailed { to, error } 报告给协议核心。一个对端地址不可达不应拖垮其他对端,所以被拒绝的发送永远不是致命的;致命的 link 错误会在读取一侧暴露。

  • Shutdown。 没有可关闭的写方向。数据包排空后,待处理的 ShutdownTransport 只是把写方向标记为已关闭。

  • 结束。 数据报传输层从不交付 TransportEof,所以这样的运行时只会通过 Effect::Finish、RuntimeError(例如 link 的致命读取错误)或被 drop 而结束。

Effect::Finish 设置 finished。由于 effect 按顺序应用,当它生效时,协议核心在 Finish 之前推入的所有内容(例如对某个出站的 Shutdown)都已经完成,或者随 key 的丢失被丢弃。

从此以后,step 不再读取传输层或任何出站。它继续应用仍在排队的内容,并把 staging 写到传输层。当 staging 为空,并且(如果请求了 ShutdownTransport)传输层的写方向已经 shutdown 时,运行时完成。quiet 形式 resolve 为 Ok(summary),stream 形式以 None 结束。传输层和仍然打开的出站在调用方 drop 已完成的运行时时关闭。

协议核心推入 结果
ShutdownTransport,然后 Finish staging 排空,传输层被 flush 并半关闭,运行时完成。客户端读到流结束。
只有 Finish staging 排空,运行时完成。运行时会 poll 一次传输层的 flush,但不等待它完成,传输层在运行时被 drop 时关闭。如果最后的字节必须被 flush,也要推入 ShutdownTransport。
在所有来源都关闭后什么都不推入 运行时保持挂起。协议核心必须推入 Finish 才能结束连接,通常是在传输层到达流结束、且所有出站都已消失之后。
concepts/src/runtime.rs
pub struct Traffic {
pub transport_rx: u64,
pub transport_tx: u64,
pub outbound_rx: u64,
pub outbound_tx: u64,
}
impl std::ops::AddAssign for Traffic { /* field-wise += */ }
计数器 计数时机
transport_rx 一次 stream 读取把 n 字节读入读缓冲区,或者收到一个 n 字节(截断后)的传输层数据包
transport_tx 一次 stream 写入接受了 n 字节,或者发送了一个传输层数据包。被拒绝的数据包不计入。
outbound_rx 一次 n 字节的出站读取或数据包落入 scratch
outbound_tx 一次转发写出 n 字节,或者一次数据报发送报告了 n 字节
concepts/src/runtime.rs
impl<const BUF_SIZE: usize, Core, Trans, Conn> Future
for ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Trans: Transport<Addr = Core::TransportAddr>,
Conn: Connector<Core::Target>,
Conn::Datagram: DatagramLink<Addr = Destination>,
{
type Output = Result<Traffic, RuntimeError<Core::Error>>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}

该 future 只 resolve 一次,结果是连接整个生命周期内的总流量。每次 poll 最多执行 WORK_BUDGET 个工作单元。Hysteria 2 运行时和 TUN UDP 运行时使用这种形式,在它们 spawn 的 task 中 await 它。测试则直接 spawn 它,例如 tokio::spawn(ProxyServerRuntime::<256, _, _, _>::new(transport, core, connector))。

stream 形式的存在是为了让调用方能在两步之间采取行动。app/src/serve.rs → drive、katana 的 src/serve.rs → drive 以及 protocols/src/tun/inbound.rs → serve_stream 中的服务循环都以同样的方式使用它:在 HANDSHAKE_TIMEOUT 限制下 poll runtime.next(),直到 runtime.core().is_established(),然后继续 poll 直到 stream 结束。沉默的客户端不会产生任何运行时事件,这就是握手看门狗放在调用方的原因。katana 还会让 next() 与它的用户下线信号竞争,并在每个非零增量时重置 PROGRESS_WATCHDOG 定时器。见入站服务和 katana 入站服务。

concepts/src/runtime.rs
pub enum RuntimeError<E> {
Transport(io::Error),
Core(E),
UnknownKey,
DuplicateKey,
RangeOutOfBounds,
BadConsume,
FrameTooLarge,
WrongLinkKind,
StagedWithoutPeer,
}
impl<E: fmt::Display> fmt::Display for RuntimeError<E>;
impl<E: fmt::Debug + fmt::Display> std::error::Error for RuntimeError<E>;
impl<E> From<io::Error> for RuntimeError<E>;

RuntimeError 会结束连接。除 Transport、Core 和 FrameTooLarge 外,每个变体都表示协议核心违反了运行时的约定,所以它指向的是某个协议核心中的 bug,而不是行为异常的对端。WrongLinkKind 也可能表示 connector 返回的 link 种类与协议核心预期的相反。From<io::Error> 映射到 Transport。

变体 Display 何时产生
Transport(e) transport: {e} 传输层的读取、写入、flush 或 shutdown 返回错误;stream 写入接受了 0 字节(WriteZero);数据报传输层返回了 stream 读取。数据报发送被拒绝不是错误。
Core(e) proxy core: {e} Core::handle 对任意事件返回 Err
UnknownKey effect targets an unknown outbound key 非空的 Forward 或 ForwardHeld、SendTo 或 SendToHeld,或者 Shutdown 到达队首时,其 key 没有槽位
DuplicateKey open reuses a live outbound key Open 到达队首时,其 key 已有槽位
RangeOutOfBounds forward range outside the event slice or held buffer 入队时,区间反向或超出事件切片或 held() 的末尾;应用时,held 区间已不再位于 held() 之内
BadConsume core consumed an impossible byte count 一次 Transport 调用报告的消费量超过所给数据,或者 TransportDatagram、Outbound 或 Datagram 的载荷没有被整体消费
FrameTooLarge protocol frame exceeds the read buffer 未解析的传输层数据占满了 BUF_SIZE,而协议核心仍需要更多数据
WrongLinkKind stream effect on a datagram outbound or vice versa 拨号完成后发现 Forward/ForwardHeld 作用于数据报槽位,或 SendTo/SendToHeld 作用于 stream 槽位
StagedWithoutPeer bytes staged toward a datagram transport without a peer 在数据报传输层下,协议核心在 stage_to/put_to 之外 stage 了字节

出站失败本身从不结束连接。它们交给协议核心,由协议核心决定:

事件 何时产生 之后的 key
ConnectFailed { key, error } 拨号 future resolve 为 Err 已移除;针对它排队的 effect 被丢弃
OutboundError { key, error } 写入返回错误或 0 字节,flush 或 shutdown 失败,stream 读取失败,或者数据报接收失败 已移除;针对它排队的 effect 被丢弃
SendFailed { key, to, error } 数据报出站拒绝了一个数据包(例如 UdpOutbound 以 Unsupported 拒绝域名) 仍存活
TransportSendFailed { to, error } 数据报传输层拒绝了一个数据包 不适用;运行时继续运行

fail_outbound 和 report_send_failure 在 drive_effects 内部运行,所以它们只把协议核心的应对入队;由产生它们的那一轮处理或下一个 step 来应用。

poll_next 以 Some(Err(e)) 产出一次错误,设置 failed,之后的每次 poll 都返回 None。

drop 这个 future 或 stream 就会取消连接。运行时持有传输层、协议核心、每个进行中的拨号 future 和每个出站 link,所以它们都随之被 drop,其 socket 也随之关闭。drop 不会 flush:仍在 staging 中的字节被丢弃,排队中的 effect 永远不会被应用。没有单独的取消方法,运行时也不会 spawn 任何可能比它活得更久的东西。调用方依赖这一点:用户下线时,katana 的 drive 从它的 select! 中返回;握手看门狗在超时时返回;两种情况下,drop 运行时就是全部的清理工作。

不变量 保证机制 固定它的测试
BUF_SIZE 在 staging 预留之外还留有空间 build 中的 assert! 无
排队中的带区间 effect 引用的字节不会被移动或覆盖 来源锁定:pinned(Source::Transport) 阻止传输层读取、压缩和交付;scratch 和 held 锁定把出站 key 暂存到 starved;held_free() 阻止字节事件 held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites、stalled_outbound_holds_uplink_but_not_other_downlink
effect 按推入顺序应用 drive_effects 从前往后处理,并在第一个被阻塞的 effect 处停下 half_close_propagates_both_ways_and_finishes、stalled_outbound_holds_uplink_but_not_other_downlink
发往仍在拨号的 key 的转发会等待拨号完成 LinkState::Connecting 使 drive_effects 跳出 held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites
协议帧能放进读缓冲区,否则连接失败 is_saturated() 检查产生 FrameTooLarge frame_larger_than_the_buffer_is_an_error
出站数据块总能被封装发往传输层 service_key 中的空间检查;stream 读取上限为 room - STAGING_RESERVE 每个中继测试都会覆盖,例如 relays_two_keys_and_completes_on_fin
数据报出站永远不会因背压而被截断 STAGING_RESERVE + datagram_limit() 的空间检查 datagram_outbound_is_never_truncated_by_staging_backpressure
发往数据报传输层的每个字节都有对端 deliver 中 stage 字节数与成帧字节数的比较 staging_without_a_peer_under_a_datagram_transport_is_an_error、stage_to_is_refused_under_a_stream_transport
一个被拒绝的数据包不会结束连接 poll_transport_send_packets 丢弃并报告;drive_effects 报告 SendFailed 并保留 key a_refused_transport_packet_is_reported_not_fatal、a_refused_datagram_send_keeps_the_key_alive、single_peer_datagram_transport_demuxes_sessions_and_reports_refusals
丢失的 key 只报告一次,之后不再向它发送任何内容 在发出那个唯一的事件之前调用 forget_key connect_failure_reaches_the_core_as_an_event 覆盖了事件;对排队 effect 的清除没有专门测试
poll 期间的唤醒不会丢失 poll 前调用 clear_queued();每次成功读取后调用 schedule() concepts/src/wake.rs 中的 clearing_queued_lets_the_key_be_queued_again
一连串出站唤醒只唤醒 task 一次 基于 queued 去重;唤醒时取走 parent waker concepts/src/wake.rs 中的 wake_queues_key_once_and_wakes_parent
deadline 不会被就绪的 I/O 饿死 在 step 中最先 poll 没有直接测试;deadline_event_lets_the_core_time_out 覆盖了触发
deadline 遵循 tokio 的时钟 set_deadline 中的 tokio::time::Instant::now() 和 sleep_until deadline_is_armed_against_the_tokio_clock
一条繁忙的连接无法独占 worker Future::poll 中的 WORK_BUDGET 循环与自我唤醒 无
增量之和等于总量 poll_once 在每个取得进展的 step 中把 delta 移入 summary stream_mode_reports_each_unit_of_work
常量 值 定义位置 含义
BUF_SIZE 按实例化而定;见运行位置 ProxyServerRuntime 的 const 泛型 三个缓冲区各自的大小,也是最大协议帧的大小
ProxyCoreDecode::STAGING_RESERVE 由各协议核心声明 concepts/src/core.rs 读取传输层之前保证可用的 staging 空间
ProxyCoreDecode::MAX_DATAGRAM 默认 4096 concepts/src/core.rs 两侧都能完整交付的最大数据报
datagram_limit() min(MAX_DATAGRAM, BUF_SIZE - STAGING_RESERVE) concepts/src/runtime.rs 数据报出站的读取大小
传输层数据包读取大小 MAX_DATAGRAM.min(BUF_SIZE) poll_transport_recv_packet 数据报传输层的读取大小
WORK_BUDGET 64 concepts/src/runtime.rs quiet 形式的每次 Future::poll 在自我唤醒之前执行的工作单元数
INLINE_EFFECTS 4 concepts/src/core.rs EffectList 和队列在溢出到堆之前内联容纳的 effect 数(一次握手是 Open + Forward + SetDeadline)

连接内存的固定部分是三个 BUF_SIZE 数组。每个打开的出站增加一个 map 条目、它的 link 和一个 Arc<KeyWaker>,连接期间还有一个 boxed 的拨号 future。定时器在第一次启用时增加一个 boxed 的 Sleep。effect 队列超过 INLINE_EFFECTS 项后溢出到堆上,而 held() 中保留多少数据由协议核心自己限定。

集成测试位于 concepts/tests/runtime.rs。它们用玩具 sans-I/O 协议核心驱动运行时,运行在内存 duplex stream、回环 UDP socket 和基于 mpsc 的数据报 link 之上:

玩具协议核心 传输层 Key STAGING_RESERVE 线上格式
TinyMux Stream u8 HDR + 64 = 68 [kind:u8][key:u8][len:u16][payload],载荷与 0x55 异或。客户端 kind:1 打开,2 数据,3 关闭,4 fin。应答:2 数据,5 已连接,6 连接失败或出站错误,7 EOF。
TinyUdp Stream Single 4 [len:u16][port:u16][payload];len == 0 表示结束
TinyHub 数据报,SocketAddr 对端 SocketAddr 2 [port:u16][payload];端口 0 表示结束
Sniffing Stream Single 0 把前 8 个字节作为目标名保留,并从 held() 转发它们
基于 TinyQuic 的 TinySessions 数据报,单对端(()) u32 6 [session:u32][port:u16][payload];会话 0 表示结束
测试 固定的行为
core_is_driven_without_any_io 协议核心在没有运行时的情况下处理手工构造的事件:被拆开的帧保持未消费,载荷被原地解密
relays_two_keys_and_completes_on_fin 两个 key 中继上行,其中一个中继下行;两个 key 的 Connected 都到达客户端;ShutdownTransport 加 Finish 完成连接;Traffic 总量精确
stalled_outbound_holds_uplink_but_not_other_downlink 部分写入会卡住它后面的上行,而另一个 key 的下行照常流动
connect_failure_reaches_the_core_as_an_event 拨号失败以 ConnectFailed 到达协议核心,而不是运行时错误
half_close_propagates_both_ways_and_finishes 待发送数据在 shutdown 之前到达出站;出站 EOF 到达协议核心;运行时结束
deadline_event_lets_the_core_time_out 在暂停的时钟下,SetDeadline 触发 Deadline
frame_larger_than_the_buffer_is_an_error BUF_SIZE = 128 时产生 FrameTooLarge
stream_mode_reports_each_unit_of_work showing_progress 产出的增量之和等于预期总量,其中一个增量携带转发,元操作 step 的增量为零
datagrams_round_trip_through_a_real_udp_socket SendTo 和 Event::Datagram 经过真实的 UDP 出站
datagram_outbound_is_never_truncated_by_staging_backpressure 客户端不读取时,1000 字节的应答在 socket 处等待,而不是被截断
deadline_is_armed_against_the_tokio_clock deadline 使用 tokio 的时钟,而不是 std 的时钟
datagram_transport_demultiplexes_peers_and_frames_replies 多个对端下的 over_datagrams;每个对端只收到自己的应答;Traffic 总量精确
staging_without_a_peer_under_a_datagram_transport_is_an_error StagedWithoutPeer
stage_to_is_refused_under_a_stream_transport stream 传输层下 stage_to 返回 None,而 stage 仍然可用
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites held 转发等待拨号完成,之后的传输层字节等待它完成
a_held_range_past_the_buffer_is_rejected held 区间越界时产生 RangeOutOfBounds
a_refused_datagram_send_keeps_the_key_alive 报告 SendFailed,同一个 key 继续中继
a_refused_transport_packet_is_reported_not_fatal 报告 TransportSendFailed,被拒绝的数据包不计入 transport_tx
single_peer_datagram_transport_demuxes_sessions_and_reports_refusals 单对端数据报 link(TransportAddr = ()):超大应答被拒绝,其他会话继续进行

代码旁边的单元测试:

文件 测试 固定的行为
concepts/src/buffer.rs read_buffer_compacts_unparsed_tail_to_front compact 把未解析的尾部移到偏移 0
concepts/src/buffer.rs read_buffer_saturation_means_frame_too_large 只有当未解析区域占满数组时 is_saturated 才成立
concepts/src/buffer.rs write_buffer_slides_pending_when_full 尾部满时 room() 滑动待写出字节;完全排空后重置
concepts/src/buffer.rs staging_refuses_over_reservation_without_partial_commit 被拒绝的 reserve 不占用任何空间
concepts/src/wake.rs wake_queues_key_once_and_wakes_parent 去重,以及每次注册只唤醒 parent 一次
concepts/src/wake.rs clearing_queued_lets_the_key_be_queued_again 只有在 clear_queued 之后 key 才能再次入队

其他运行服务端运行时的测试:

  • concepts/tests/client.rs 的 server_side 模块:服务端运行时的出站是同一 task 中的一个客户端运行时,见 server_runtime_relays_through_a_client_runtime_in_one_task、upstream_dial_failure_is_connect_failed_not_connected 和 refused_upstream_handshake_is_connect_failed_not_connected。
  • environment/tests/integration/udp.rs → dual_stack_link_serves_a_proxy_runtime:运行在 environment crate 双栈 UDP link 之上的运行时。
  • protocols/tests/support/pipeline.rs → serve_runtime:协议测试使用的测试框架,为每个 accept 的回环连接运行一个运行时。

变体 UnknownKey、DuplicateKey、WrongLinkKind 和 BadConsume,以及 WORK_BUDGET 自我唤醒,都没有专门的测试。修改这些路径时请补上测试。不经运行时测试协议核心的方法见测试。