跳转到内容

服务端协议核心

源码文件:33 个 · 核对版本 Etemenanki 596916d
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/tests/runtime.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/core/harness.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/tests/unit/mux/demux.rs
  • Etemenanki/protocols/src/ss_legacy/aead.rs
  • Etemenanki/protocols/src/vmess/framing.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/sniff/collector.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/http/protocol.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/tests/unit/ss_legacy/aead.rs
  • Etemenanki/protocols/src/ss_2022/core.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/tests/unit/vmess/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/udp.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/tests/unit/socks/server.rs
  • Etemenanki/app/src/serve.rs

Etemenanki 中除 SOCKS 之外的每个服务端协议都是一个协议核心(core):一个实现了 ProxyCoreDecode 的普通 Rust 值,从不接触 socket、时钟或任务 Context。I/O 由每连接运行时负责:它读取字节,每次把一个 Event 交给协议核心,并执行协议核心推回的 Effect。由于协议核心是一个从字节和事件映射到字节和 effect 的函数,测试可以用手工构造的输入驱动它,并对输出做断言。SOCKS 是例外,因为它的 UDP 部分位于第二个 socket 上,即一个中继“hub”,控制连接只负责让它保持存活;protocols/src/socks/server.rs 在这份契约之外把它作为单个任务运行。这个任务还决定 hub 接收谁的数据报:ExpectedSender 把每个 UDP 关联限定在控制连接的客户端上,要求每个数据报都来自对端的 IP(由 protocols/src/socks/protocol.rs 中的 endpoint 按规范形式比较),并以它转发的第一个数据报的端口固定端口;如果 UDP ASSOCIATE 请求在同一 IP 之外还给出了非零端口,则固定为该端口。通过 Unix socket 连接时没有对端,请求必须给出确切的源地址和端口。来自其他任何地址的数据报都会被直接丢弃,不做解析。protocols/tests/unit/socks/server.rs 中的 only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port 固定了这条规则,完整的对照表见 SOCKS。

本页是协议核心与运行时之间的契约:trait、每个事件和 effect、Effects 汇集器,以及确保原地解密、held 缓冲区、key、截止时间和数据报安全的规则。编写新的协议核心、或修改 concepts/src/runtime.rs 中的投递逻辑之前,请先读本页。运行时的调度、waker 和背压见服务端运行时;客户端一侧(ProxyCoreEncode)见客户端运行时。

协议核心做决定,运行时去执行。

协议核心 运行时(ProxyServerRuntime)
解析客户端的线上字节,原地解密 把传输层数据读入固定大小的读缓冲区
决定每个流的去向,并为它生成一个 key(Effect::Open) 通过自己的 Connector 拨号到目标,并为每个 key 维护一个槽位
指明哪些字节区间发往哪个出站(Effect::Forward) 带背压地按顺序写出这些区间
把回复字节加密封装进 staging 区 把 staging 区写入传输层
自己维护定时器,并设置其中最早的一个(Effect::SetDeadline) 持有唯一的 tokio 定时器,并投递 Event::Deadline
宣告连接结束(Effect::Finish) 排空、关闭,并以 Traffic 总量作为结果返回

随机数、挂钟时间和共享的账户状态都是具体协议核心的构造参数,从不由运行时注入。这使得对给定的协议核心值,每次调用都是确定性的。VMessCore::new 接收 now: fn() -> i64,Ss2022Core::new 接收 now: fn() -> u64(生产环境使用的构造函数是 Ss2022Core::with_system_clock),Hy2UdpCore 则以自己的清扫计数作为时钟。

整份契约都位于 concepts/src/core.rs。树内协议核心所组合的共享构件位于 protocols/src/core/mod.rs。

concepts/src/core.rs
pub struct Single;
pub trait ProxyCoreDecode {
type Key: Copy + Ord + Send + Sync + 'static;
type Target;
type Error;
type TransportAddr: Clone + Send + Sync + 'static;
const STAGING_RESERVE: usize;
const MAX_DATAGRAM: usize = 4096;
fn handle(
&mut self,
event: Event<'_, Self>,
effects: &mut Effects<'_, Self>,
) -> Result<usize, Self::Error>;
fn held(&self) -> &[u8] {
&[]
}
}
项 含义
Key 标识一个出站。协议核心通过 Effect::Open 生成 key,每个出站侧事件都带有其来源的 key。只有一个出站的协议核心使用单元结构体 Single。运行时把槽位存放在 BTreeMap 中,因此要求 Ord。
Target 运行时的 Connector 拨号的对象。对树内服务端协议核心而言是 Flow<T>(目的地、用户、嗅探结果、源地址)。
Error 由 handle 返回;运行时将其包装为 RuntimeError::Core 并结束连接。树内协议核心使用 io::Error。
TransportAddr 数据报传输层中数据包的对端地址:Event::TransportDatagram 的来源,也是 Effects::stage_to 的发送目标。一个 UDP socket 服务多个对端时为 SocketAddr;QUIC 连接的数据报(单一对端)以及任何字节流协议核心为 ()。
STAGING_RESERVE 协议核心在处理单个事件时,除回显该事件自身负载之外,向传输层 staging 的字节数上限。见 staging 预留。
MAX_DATAGRAM 能完整投递的最大数据报,适用于数据报出站和数据报传输层。更长的数据包会像内核 recv 那样被截断。默认为 4096。
handle 让状态机前进一个事件。返回携带字节的事件中,从头部起被消费的字节数。
held 协议核心自行保管、并通过 ForwardHeld 或 SendToHeld 按区间转发的字节。默认为空。
concepts/src/core.rs
pub enum Event<'a, C: ProxyCoreDecode + ?Sized> {
Transport(&'a mut [u8]),
TransportDatagram {
from: C::TransportAddr,
data: &'a mut [u8],
},
TransportSendFailed {
to: C::TransportAddr,
error: io::Error,
},
TransportEof,
Outbound { key: C::Key, data: &'a mut [u8] },
Datagram {
key: C::Key,
from: Destination,
data: &'a mut [u8],
},
SendFailed {
key: C::Key,
to: Destination,
error: io::Error,
},
OutboundEof { key: C::Key },
Connected { key: C::Key },
ConnectFailed { key: C::Key, error: io::Error },
OutboundError { key: C::Key, error: io::Error },
Deadline,
}

Debug 实现只打印切片长度,从不打印切片内容。

事件 运行时何时投递 handle 必须返回什么
Transport(data) 字节流传输层产生了字节。data 是读缓冲区中完整的未解析区域,包括协议核心上次调用留下的尾部。 从头部消费的字节数,范围 0..=data.len()。剩余部分会再次提供。
TransportDatagram { from, data } 数据报传输层产生了一个完整的数据包。在流式传输层上从不投递。 恰好 data.len()。
TransportSendFailed { to, error } 发往对端 to 的一个已 staging 的数据包被链路拒绝并丢弃。不致命。 按惯例返回 0;不检查该计数。
TransportEof 流式传输层的读端结束。数据报传输层从不投递它。 按惯例返回 0。
Outbound { key, data } 流式出站 key 产生了字节(读入 scratch 缓冲区)。 恰好 data.len():没有按出站划分的缓冲区可以保留剩余部分。
Datagram { key, from, data } 数据报出站 key 产生了一个数据包。from 是数据包的来源;会解析域名的链路报告的是实际收到数据的地址,形式为 IP。 恰好 data.len()。
SendFailed { key, to, error } 数据报出站 key 上的一个 Effect::SendTo 或 SendToHeld 被拒绝。该数据包被丢弃,key 保持存活。 按惯例返回 0。
OutboundEof { key } 出站 key 的读端结束。 按惯例返回 0。
Connected { key } 由 Effect::Open 为 key 发起的拨号已完成。 按惯例返回 0。
ConnectFailed { key, error } 拨号失败。该 key 已经不存在。 按惯例返回 0。
OutboundError { key, error } 出站 key 在读、写、flush 或关闭时失败。该 key 已经不存在。 按惯例返回 0。
Deadline 由 Effect::SetDeadline 设置的截止时间已过,定时器现已解除。 按惯例返回 0。

运行时在 feed_transport、poll_transport_recv_packet 和 after_scratch 中强制执行这些消费规则:消费量超过切片长度,或少于一个完整的出站负载或传输层数据报,都会以 RuntimeError::BadConsume 结束连接。

concepts/src/core.rs
pub enum Effect<C: ProxyCoreDecode + ?Sized> {
Forward { key: C::Key, range: Range<usize> },
SendTo {
key: C::Key,
to: Destination,
range: Range<usize>,
},
ForwardHeld { key: C::Key, range: Range<usize> },
SendToHeld {
key: C::Key,
to: Destination,
range: Range<usize>,
},
Open { key: C::Key, target: C::Target },
Shutdown { key: C::Key },
Close { key: C::Key },
ShutdownTransport,
SetDeadline(Option<Duration>),
Finish,
}

effect 由 ProxyServerRuntime::drive_effects 严格按推入顺序应用。暂时无法完成的 effect 会留在队首,其后的所有 effect 都要等待,不论它们指向哪个 key。在此期间,运行时让该事件的字节保持有效,也不会向那个缓冲区读入新数据。

Effect 应用时做什么 会阻塞队列吗?
Forward 把事件切片中的 range 写入流式出站 key,然后 flush。短写会推进 range.start 并立即再次写入,直到写操作返回 Pending。写入零字节会以 WriteZero 使该出站失败。空区间会被丢弃。 会:当 key 仍在连接中,或写操作返回 Pending 时。
SendTo 通过数据报出站 key,把事件切片中的 range 作为一个数据报发往 to。 会:连接中,或发送返回 Pending 时。
ForwardHeld 同 Forward,但 range 索引的是 held(),并在写入的那一刻解析。 会,同 Forward。
SendToHeld 同 SendTo,但 range 索引的是 held()。 会,同 SendTo。
Open 调用 connector.connect(target),把拨号 future 存入为 key 新建的槽位,并调度该 key。 不会。key 被调度时,运行时会与该 key 的其他工作一并轮询这个拨号 future;排在它之后发往 key 的转发会等待连接完成。排在某个阻塞 effect 之后的 Open 不会被应用,其拨号也不会开始,直到那个 effect 完成。
Shutdown 在其之前排队的所有数据都写出后,半关闭流式出站 key 的写端。对数据报出站则立即完成。 会:连接中,或 poll_shutdown 返回 Pending 时。
Close 双向丢弃 key 的槽位,包括仍在进行的拨号。之后不会有任何事件。 不会。
ShutdownTransport 设置一个标志。staging 排空并 flush 后,运行时会半关闭传输层的写端。 不会。
SetDeadline(after) Some(d) 把唯一的定时器设为从现在起 d 之后触发,取代之前的任何截止时间。None 解除定时器。 不会。
Finish 标记连接结束。运行时停止读取,写出 staging,执行待处理的 ShutdownTransport,然后返回结果。 不会。
concepts/src/core.rs
pub const INLINE_EFFECTS: usize = 4;
pub type EffectList<C> = SmallVec<[Effect<C>; INLINE_EFFECTS]>;
pub type PacketList<A> = VecDeque<(usize, A)>;
pub struct Effects<'a, C: ProxyCoreDecode + ?Sized> { /* private */ }
impl<'a, C: ProxyCoreDecode + ?Sized> Effects<'a, C> {
pub fn new(list: &'a mut EffectList<C>, staging: Staging<'a>) -> Self;
pub fn with_packets(
list: &'a mut EffectList<C>,
staging: Staging<'a>,
packets: &'a mut PacketList<C::TransportAddr>,
) -> Self;
pub fn push(&mut self, effect: Effect<C>);
pub fn staging(&mut self) -> &mut Staging<'a>;
pub fn stage(&mut self, bytes: &[u8]) -> Option<()>;
pub fn transport_is_datagram(&self) -> bool;
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 list(&self) -> &[Effect<C>];
}
方法 用途
push 追加一个 effect。列表归运行时所有,跨调用复用,最多 INLINE_EFFECTS 个时内联存储(一次握手通常是 Open + Forward + SetDeadline)。
staging 发往传输层的 Staging 区。Staging::reserve(len) 在尾部占用 len 字节并返回供填充,帧就是这样原地加密封装的,无需中间 Vec。Staging::room 给出剩余空间。
stage 把字节复制到发往字节流传输层的方向。空间不足时返回 None。
transport_is_datagram 汇集器通过 with_packets 构建时为 true,此时 staging 的字节需要一个对端。
stage_to 占用 len 字节作为发往 to 的一个数据包,并返回供填充。空间不足或传输层是字节流时返回 None。
put_to 把 bytes 复制为发往 to 的一个数据包。
list 到目前为止推入的 effect,供协议核心在测试中检查自己的输出。

Staging 是汇集器所包装的、运行时发往传输层的 WriteBuffer<BUF_SIZE> 的只追加视图:

concepts/src/buffer.rs
pub struct Staging<'a> { /* private */ }
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<()>;
}

reserve 在尾部占用 len 字节,并立即将其计为已 staging;协议核心原地填充它们。put 就是 reserve 加一次复制,Effects::stage 即 put。staging 从不部分提交:放不下的 reserve 返回 None,且什么都不会 staging(见 concepts/src/buffer.rs 中的 staging_refuses_over_reservation_without_partial_commit)。在交出 Staging 之前,WriteBuffer::staging 会调用 WriteBuffer::room,后者在空闲尾部用完时把仍待写出的字节滑到最前面;这次滑动是 staging 缓冲区自行进行的唯一一次复制(write_buffer_slides_pending_when_full)。

用 Effects::new 构建。用 stage 或 staging().reserve(..) 追加。运行时把 WriteBuffer 中全部待写区域写入传输层。

一条连接对应一个运行时和一个协议核心。运行时在自己的 poll 内部同步调用 handle,因此协议核心从不与自身并发运行,也不需要锁。

sequenceDiagram
  participant T as Transport
  participant R as ProxyServerRuntime
  participant C as Core
  participant O as Outbound
  T->>R: 字节读入读缓冲区
  R->>C: handle(Event::Transport(unparsed))
  C-->>R: 消费计数,Open(k) + Forward(k, range)
  R->>R: 校验区间,按序入队
  R->>O: 应用 Open:connector.connect(target)
  Note over R: k 连接期间 Forward 等待
  O-->>R: 拨号完成
  R->>C: handle(Event::Connected(k))
  R->>O: 应用 Forward:写出区间
  O-->>R: 回复字节读入 scratch
  R->>C: handle(Event::Outbound(k, data))
  C-->>R: 加密封装后的回复写入 staging
  R->>T: 写出 staging

无论是什么事件,每次投递都经过相同的三步。ProxyServerRuntime::deliver 构建汇集器并调用 handle;enqueue 校验每个带区间的 effect,把基于事件切片的区间换算到它所索引的缓冲区上,并把 effect 追加到队列;随后 drive_effects 从队首开始应用队列,直到某个 effect 被阻塞。响应通知类事件(Connected、Deadline、OutboundError 等)时推入的 effect 会加入同一队列的末尾。

flowchart TB
  ev["事件就绪且允许投递"] --> deliver["deliver: core.handle(event, effects)"]
  deliver --> check{"消费计数有效?"}
  check -- 否 --> bad["RuntimeError::BadConsume"]
  check -- 是 --> enq["enqueue:检查区间、换算、标记来源缓冲区"]
  enq --> oob{"区间在界内?"}
  oob -- 否 --> rng["RuntimeError::RangeOutOfBounds"]
  oob -- 是 --> drive["drive_effects:应用队首 effect"]
  drive --> blocked{"队首 effect 被阻塞?"}
  blocked -- "否,还有排队" --> drive
  blocked -- "是,或队列为空" --> wait["返回调度器"]

concepts/tests/runtime.rs 驱动一个玩具级多路复用协议 TinyMux,其帧格式为 [kind][key][len:u16][payload],并使用 XOR“加密”(XOR = 0x55,HDR = 4)。它声明 Key = u8、Target = String、TransportAddr = () 和 STAGING_RESERVE = HDR + 64。它处理传输层事件的分支就是标准的解析循环:

concepts/tests/runtime.rs (excerpt)
Event::Transport(buf) => {
let mut consumed = 0;
loop {
let rest = &mut buf[consumed..];
if rest.len() < HDR {
break;
}
let (k, key) = (rest[0], rest[1]);
let len = u16::from_be_bytes([rest[2], rest[3]]) as usize;
if rest.len() < HDR + len {
break;
}
let payload = &mut rest[HDR..HDR + len];
for b in payload.iter_mut() {
*b ^= XOR;
}
let (start, end) = (consumed + HDR, consumed + HDR + len);
match k {
kind::OPEN => { /* 插入 key,推入 Effect::Open */ }
kind::DATA => fx.push(Effect::Forward { key, range: start..end }),
kind::CLOSE => fx.push(Effect::Shutdown { key }),
kind::FIN => { /* ShutdownTransport、Finish */ }
other => return Err(format!("bad frame kind {other}")),
}
consumed = end;
}
/* 可选的空闲 SetDeadline */
Ok(consumed)
}

core_is_driven_without_any_io 完全不借助运行时来驱动它:

  1. 输入。 测试构造 OPEN 9 "somewhere"、DATA 9 "hi" 以及 DATA 9 "split" 的前 5 个字节,并在一个基于 WriteBuffer::<256> 的 Effects::new 汇集器上调用 handle(Event::Transport(&mut input), &mut fx)。

  2. 逐帧解析。 协议核心解码两个完整的帧,并在不完整的第三帧处停下。它返回 HDR + 9 + HDR + 2,即被截断的帧开始的偏移。运行时会保留这 5 个字节,并在下次读取时把它们放在新数据前面再次提供。

  3. Effect。 列表中依次是 Open { key: 9, target: "somewhere" } 和 Forward { key: 9, range: 17..19 }。该区间相对于协议核心拿到的切片;运行时在把 effect 入队时将其换算到读缓冲区上。

  4. 原地处理。 input[17..19] 现在是 hi:负载在原位置被解密,转发时无需复制即可写出。

  5. Connected。 handle(Event::Connected { key: 9 }) 把帧 [5, 9, 0, 0](kind::CONNECTED)写入 staging。不需要任何 effect。

  6. 下行。 handle(Event::Outbound { key: 9, data: b"pong" }) 用 fx.staging().reserve(HDR + len) 把一个 DATA 帧加密封装进 staging,并返回 4,即整个负载。

随后同一个协议核心在真实的 ProxyServerRuntime 下、基于 tokio::io::duplex 管道运行:relays_two_keys_and_completes_on_fin 打开两个 key 并检查 Traffic 总量;stalled_outbound_holds_uplink_but_not_other_downlink 给 key 1 一个 8 字节的管道,检查在 key 1 的转发之后排队的上行等待期间,key 2 的下行仍能送达。

Event::Transport 把 ReadBuffer 中完整的未解析区域 [start..end) 交给协议核心。协议核心返回从头部消费了多少字节。运行时(feed_transport)按该计数推进 start,并在协议核心有进展、仍有剩余字节、且没有排队的 effect 占用读缓冲区或 held 缓冲区时再次调用。当未解析区域为空、或某次调用什么都没消费时,它才从传输层读取更多数据。请逐帧解析,在第一个不完整的帧处停下,并返回它开始处的偏移。

  • 机制: ReadBuffer 把已解析的字节 [0..start) 留在原位,因为排队中的转发可能仍然引用它们。当排队中的转发占用该缓冲区时,运行时不会向其中读入任何数据;只有空闲尾部用完时,才调用 compact 把未解析的尾部移到最前面。
  • 测试: concepts/tests/runtime.rs 中的 core_is_driven_without_any_io(被截断的帧保持未消费);concepts/src/buffer.rs 中的 read_buffer_compacts_unparsed_tail_to_front。

切片是 &mut,所以帧可以原地解密并按区间转发。其后果是:如果一个帧的头部已经原地解密,而帧体尚未到达,下一次调用时同样的字节会以已变换的形式再次出现,天真的解析器会把它们解密两次。请把解码出的长度,以及任何已推进的密码状态(nonce 计数器、SHAKE 掩码),保存在协议核心中,并在下一次调用时跳过头部。

  • 机制: protocols/src/ss_legacy/aead.rs 中的 ChunkDecoder 从副本中解出长度,递增自己的 NonceCounter,并保存 pending_len: Option<usize>。只有整个分块都到达后,才原地解开负载。protocols/src/vmess/framing.rs 为 VMess 另有自己的 ChunkDecoder。
  • 测试: protocols/tests/unit/vmess/core.rs 中的 a_chunk_split_across_reads_decodes_its_header_once;protocols/tests/unit/ss_legacy/aead.rs 中的 chunk_encoder_and_decoder_agree_and_never_decrypt_twice(同一个头部被提供三次,每次附带更多帧体)。

BUF_SIZE 是 ProxyServerRuntime 的 const 泛型参数,决定了它三个缓冲区各自的大小:传输层读缓冲区、传输层 staging 缓冲区和出站 scratch 缓冲区。当未解析区域占满整个读缓冲区、而协议核心仍然什么都不消费时,运行时以 RuntimeError::FrameTooLarge 失败。请把 BUF_SIZE 设为协议允许的最大帧:一个 AEAD 分块加上其开销、一个完整的 HTTP 头部、一个带元数据的 mux 帧。每个树内协议核心都以关联常量 BUF_SIZE 公布其运行时所需的值(见限制)。

  • 机制: 当读缓冲区没有空闲尾部时,poll_transport_read 在调用 compact 之前先检查 ReadBuffer::is_saturated(start == 0 && end == N)。饱和的缓冲区是压缩也无法腾出空间的缓冲区,因此运行时返回 FrameTooLarge。
  • 测试: concepts/tests/runtime.rs 中的 frame_larger_than_the_buffer_is_an_error(把 500 字节的帧送入 BUF_SIZE = 128);concepts/src/buffer.rs 中的 read_buffer_saturation_means_frame_too_large。

ProxyServerRuntime::build 还会断言 BUF_SIZE > Core::STAGING_RESERVE,否则 panic 并给出 “BUF_SIZE must exceed the core’s STAGING_RESERVE or nothing is ever read”。

协议核心可以转发不在事件切片中的字节:即它自行保管、并通过 held() 暴露的字节。代码树中的用法:

  • 嗅探。 把流的前若干字节复制进 held 缓冲区并消费掉,设置 SetDeadline,然后持续收集,直到找到 TLS SNI 或 HTTP Host、预算用完或截止时间到达。随后推入带嗅探目标的 Open,再推入一个转发全部已收集数据的 ForwardHeld。SniffPrefix 为每个嗅探连接自身流的树内协议核心完成这件事。mux.cool 子流是例外:Demux 只嗅探 New 帧的负载并立即打开,因为为一个子流扣住承载连接会让其他子流停滞。
  • 改写后的请求,例如普通 HTTP 代理转发的请求头,它必须送达出站,但从未在线上出现过。
  • 跨分块重组的帧,例如跨越多个 VMess 分块的 mux 帧(Demux::feed_chunks),或由分片重组而成的 Hysteria UDP 数据包(Hy2UdpCore 用 SendToHeld 发送它)。Demux::feed_chunks 在一次调用中接收同一个字节事件解开的所有分块:它推送的每个 held 转发都在事件结束后才读取 held 缓冲区,因此该事件中较早完成的帧必须原地保留到那时。

保证其安全的规则是:held 区间在其 effect 被应用时解析,而不是在推入时解析。如果出站仍在连接或不可写,这可能发生在几次调用之后。在所有排队的 held effect 都被应用之前,运行时不会投递任何字节事件,因此在任何字节事件开始时,协议核心都可以随意清空或覆盖 held 缓冲区。在此期间仍可能到达的通知类事件只能向其追加;追加时会重新分配的 Vec 也没有问题,因为每次尝试写入时区间都会针对当前切片重新解析。

事件 触发位置 事先检查的 staging 空间 当排队的 effect 占用以下缓冲区时暂缓投递
Transport poll_transport_read、feed_transport STAGING_RESERVE,每轮读取开始时检查一次 读缓冲区或 held 缓冲区
TransportDatagram poll_transport_recv_packet STAGING_RESERVE 读缓冲区或 held 缓冲区
TransportEof poll_transport_read STAGING_RESERVE 读缓冲区或 held 缓冲区
Outbound、OutboundEof service_key STAGING_RESERVE + 1 scratch 缓冲区或 held 缓冲区
Datagram service_key STAGING_RESERVE + datagram_limit() scratch 缓冲区或 held 缓冲区
Connected、ConnectFailed service_key,在 key 连接期间 STAGING_RESERVE + 1 从不暂缓
OutboundError service_key(读或 flush 错误)或 drive_effects(写或关闭错误) 仅针对读错误,此时需要先通过该次读取自身的空间检查 读错误只会在读取时出现,而读取会像 Outbound 一样被暂缓;flush、写和关闭错误从不暂缓
SendFailed drive_effects 不检查 从不暂缓
TransportSendFailed poll_transport_send_packets 不检查 从不暂缓
Deadline poll_deadline 不检查 从不暂缓

held 缓冲区的大小由协议核心自行决定,它也是运行时唯一不设上限的每连接内存分配。请按协议允许的范围对它设限;SniffPrefix 以 SNIFF_LIMIT 为上限。

  • 机制: Queued::pins(Source::Held) 和 ProxyServerRuntime::held_free,在 feed_transport、poll_transport_read 和 service_key 中检查。held 区间在入队时(enqueue)和应用时(drive_effects)各针对 held().len() 校验一次;越界即为 RuntimeError::RangeOutOfBounds。
  • 测试: held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites(一个受控拨号让 ForwardHeld 保持排队,同时有更多客户端字节到达;这些字节只在它之后到达出站)和 a_held_range_past_the_buffer_is_rejected,都在 concepts/tests/runtime.rs 中;protocols/tests/unit/mux/demux.rs 中的 vmess_keeps_every_frame_one_read_completes(一次读取解开三个 VMess 分块,其中两个帧跨越分块边界,两者都从 held 缓冲区转发)。

STAGING_RESERVE 是协议核心在一次调用中,除所给负载之外最多 staging 的字节数:一个回复头部,或单帧开销乘以一次出站读取可能被拆分成的帧数。只有当 staging 有 STAGING_RESERVE + 1 字节空闲时,运行时才轮询流式出站,并且最多向 scratch 读入 room - STAGING_RESERVE 字节(上限为 BUF_SIZE)。因此一个 n 字节的 Outbound 事件总是伴随至少 STAGING_RESERVE + n 字节的空间,加密封装它永远不会失败,Staging::reserve 可以直接 unwrap,或者像 Passthrough::on_outbound 那样映射为 staging_full() 错误。

传输层事件只在读取之前检查一次预留;此后只要协议核心有进展,feed_transport 就会在未解析的尾部上再次调用它,而不再检查。传输层字节预期是被转发的,而不是被 staging 的。会把传输层字节回显回去的协议核心(例如应答探测的服务端)应读取 Staging::room,如果无法 staging 自己的应答,就少消费一些字节。

Deadline 不论空间如何都会投递,这样即使对端已停止读取,空闲超时也仍能触发。在应用 effect 或发送数据包期间触发的通知类事件也是如此(见上表)。协议核心在这些处理分支中必须容忍 stage 失败。

  • 机制: service_key 中的 room_short 检查,以及 poll_transport_read 中的 staging.room() < Core::STAGING_RESERVE 检查。
  • 测试: concepts/tests/runtime.rs 中的 datagram_outbound_is_never_truncated_by_staging_backpressure;concepts/src/buffer.rs 中的 write_buffer_slides_pending_when_full。

key 从 Effect::Open 开始存活,直到其出站不复存在。运行时在以下任一情况下移除槽位:

stateDiagram-v2
  [*] --> Connecting: 应用 Effect Open
  Connecting --> Live: 拨号成功,Event Connected
  Connecting --> [*]: 拨号失败,Event ConnectFailed
  Live --> WriteClosed: 应用 Effect Shutdown
  Live --> ReadClosed: Event OutboundEof
  WriteClosed --> [*]: Event OutboundEof
  ReadClosed --> [*]: 应用 Effect Shutdown
  Live --> [*]: Effect Close,或 Event OutboundError

Close 在 key 连接期间同样有效,并会丢弃拨号。连接完成后,任何状态下的 I/O 错误都会移除该 key 并投递 OutboundError。

  • 在存活的 key 上 Open 会以 RuntimeError::DuplicateKey 使连接失败(在 drive_effects 应用该 Open 时检查)。
  • 指向未知 key 的 effect。 指向没有槽位的 key 的 Forward、SendTo、它们的 held 变体以及 Shutdown,会以 RuntimeError::UnknownKey 失败。区间为空的 Forward 或 ForwardHeld 在查找 key 之前就被丢弃;对未知 key 的 Close 什么也不做。
  • 链路类型不匹配。 在数据报出站上执行流式 effect,或者反过来,会以 RuntimeError::WrongLinkKind 失败。
  • key 因失败而消失后,forget_key 会丢弃所有仍指向它的排队 effect(等待一个刚失败的拨号的转发已无处可去),并通过 ConnectFailed 或 OutboundError 通知协议核心一次。
  • 拒绝发送的数据报出站会保留其 key;协议核心收到 SendFailed。只有接收失败才会移除数据报 key,因为那意味着 socket 或连接本身已不存在。
  • 数据报出站从不报告流结束。 对它的 Shutdown 立即完成并把写端标记为关闭,但之后不会有 OutboundEof,因此数据报 key 会一直存活,直到 Close、一次接收失败或连接结束。Hy2UdpCore 在清扫截止时间到达时用 Close 回收安静的会话;TunUdpCore 每个运行时只服务一个关联,因此改用 Finish 结束整个运行时。
  • 由协议核心自己移除的 key(Close,或第二次半关闭)会直接消失,不产生任何事件。仍指向它的排队 effect 不会被丢弃,因此之后再为它推入 Forward 或 Shutdown 会以 UnknownKey 结束连接。

如果协议的对端可能在关闭某个会话 id 后立即复用它,协议核心要么在重新打开前先 Close 旧 key,要么给 key 打上 generation 标记。mux.cool 解复用器采用后者。它的 key 是 FlowKey::Sub(SubKey):

protocols/src/core/mod.rs
pub enum FlowKey {
Direct,
Sub(SubKey),
}
pub struct SubKey {
pub id: u16,
pub generation: u32,
}

Demux 为每个承载连接维护一个 generation 计数器,每接受一个 New 帧就递增一次(wrapping_add(1))。当对端结束会话 1 并立即打开新的会话 1 时,新出站会得到一个全新的 SubKey,而旧出站可能仍在排空其读端;仍为已退役 generation 到达的下行字节不会被封装给任何人。

  • 测试: protocols/tests/unit/mux/demux.rs 中的 stream_sessions_open_forward_and_end_with_fresh_generations;concepts/tests/runtime.rs 中的 connect_failure_reaches_the_core_as_an_event 和 a_refused_datagram_send_keeps_the_key_alive。没有测试覆盖 DuplicateKey、UnknownKey、WrongLinkKind 或 BadConsume。

每条连接恰好只有一个定时器。有多个超时的协议(握手时限、每会话的空闲时限、keep-alive)自行维护一个有序的到期时间表,用 SetDeadline 设置其中最早的一个,并在每次 Deadline 时重新设置。Deadline 事件意味着定时器已触发且现已解除;再次设置会取代之前的任何截止时间。

运行时基于 tokio::time::Instant 而非 std::time::Instant 设置定时器,因此测试中暂停的 tokio 时钟会被遵守。协议核心绝不能在 handle 内调用 Instant::now();挂钟时间来自传给其构造函数的时钟,以便测试驱动。

  • 机制: ProxyServerRuntime::set_deadline(tokio::time::sleep_until,或对已有定时器调用 Sleep::reset)和 poll_deadline;step 在任何 effect 或 I/O 工作之前先轮询后者,因此一条始终 I/O 就绪的连接不会饿死自己的截止时间。
  • 测试: concepts/tests/runtime.rs 中的 deadline_event_lets_the_core_time_out 和 deadline_is_armed_against_the_tokio_clock(在运行时启动前把时钟推进一小时;5 s 的空闲截止时间在 4 s 内不得触发)。

数据报在两侧都遵循 socket 的规则:每次读取向空缓冲区读入一个数据包,整包消费;超过上限的数据包像内核 recv 那样被截断。

数据报出站 数据报传输层
以何种形式到达 Event::Datagram { key, from, data } Event::TransportDatagram { from, data }
发送方式 Effect::SendTo / SendToHeld Effects::stage_to / put_to
读取上限 datagram_limit() = min(MAX_DATAGRAM, BUF_SIZE - STAGING_RESERVE) min(MAX_DATAGRAM, BUF_SIZE)
仅当 staging 空闲达到以下值时读取 STAGING_RESERVE + datagram_limit() STAGING_RESERVE
发送被拒绝时 丢弃数据包,保留 key,投递 Event::SendFailed 丢弃数据包,运行时继续,投递 Event::TransportSendFailed
寻址方式 Destination,可以是域名;由 DatagramLink 决定如何到达 TransportAddr,即传输层链路所用的寻址方式

由于数据报出站只在有空间容纳一个完整数据包加预留时才被读取,停止读取的客户端会让下行停在出站 socket 处,而不会丢失数据包尾部。对于数据包超过 MTU 的协议,请调大 MAX_DATAGRAM,并相应调整 BUF_SIZE。

数据报传输层在一个运行时上服务多个对端,因此其协议核心是一个解复用器:从数据包的对端推导出 Key,首次见到时 Open,并通过截止时间让空闲对端过期。这一侧没有流结束;Effect::Finish 是这种运行时结束的唯一方式。被传输层拒绝的数据包会被丢弃而不是使连接失败,因为一个对端的地址不可达不能拖垮其他对端。

  • 测试: datagrams_round_trip_through_a_real_udp_socket、datagram_transport_demultiplexes_peers_and_frames_replies、staging_without_a_peer_under_a_datagram_transport_is_an_error、stage_to_is_refused_under_a_stream_transport、a_refused_transport_packet_is_reported_not_fatal 和 single_peer_datagram_transport_demuxes_sessions_and_reports_refusals,都在 concepts/tests/runtime.rs 中。

树内协议核心不继承任何基础类型。它们组合 protocols/src/core/mod.rs 中的三个小型辅助组件,对用户负载 T 泛型(应用使用 ()),以 io::Error 失败,并暴露 is_established(),让应用能区分从未发言的客户端和正在被服务的客户端。

Timing:单一截止时间,按阶段设置

Section titled “Timing:单一截止时间,按阶段设置”
protocols/src/core/mod.rs
pub enum Phase { Handshake, Sniff, Relay, Closing }
pub enum Expired { Handshake, Sniff, Idle }
pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
pub const RELAY_IDLE_TIMEOUT: Duration = Duration::from_secs(300);
impl Timing {
pub fn new() -> Self;
pub fn phase(&self) -> Phase;
pub fn is_established(&self) -> bool;
pub fn touch<C: ProxyCoreDecode>(&mut self, fx: &mut Effects<'_, C>);
pub fn enter<C: ProxyCoreDecode>(&mut self, phase: Phase, fx: &mut Effects<'_, C>);
pub fn expired<C: ProxyCoreDecode>(&mut self, fx: &mut Effects<'_, C>) -> Expired;
}
stateDiagram-v2
  [*] --> Handshake
  Handshake --> Sniff: enter Sniff
  Handshake --> Relay: enter Relay
  Sniff --> Relay: enter Relay
  Relay --> Closing: 空闲截止时间到达,推入 Finish
阶段 设置的截止时间 收到 Deadline 时 expired 返回
Handshake HANDSHAKE_TIMEOUT(10 s),由第一次 touch 设置一次 Expired::Handshake;协议核心返回 handshake_timed_out(),即一个 io::ErrorKind::TimedOut 错误
Sniff SNIFF_TIMEOUT(300 ms),由 enter(Phase::Sniff) 设置 Expired::Sniff;协议核心用已收集的数据打开流
Relay、Closing RELAY_IDLE_TIMEOUT(300 s),每个字节事件都由 touch 重新设置 Expired::Idle;Timing 已推入 Effect::Finish 并进入 Closing

正是每个字节事件都重新设置,才使中继时限成为空闲时限而非生命周期上限。Timing::new 什么都不设置:客户端发言之前运行时不会投递任何事件,所以沉默的客户端需要由应用的看门狗来处理。例如,app/src/serve.rs 中的 drive 和 protocols/src/tun/inbound.rs 中的 serve_stream 都以 showing_progress 流模式运行运行时,并在 is_established() 变为 true 之前,把每一步都包在 tokio::time::timeout(HANDSHAKE_TIMEOUT, runtime.next()) 中。

SniffPrefix 包装了嗅探用的 Collector。push(plain) 最多取用剩余的 SNIFF_LIMIT(4 KiB)预算,返回 (taken, Verdict),其中 Verdict 为 Found、More 或 Exhausted。协议核心恰好消费 taken 字节。held() 暴露已收集的字节供 ForwardHeld 使用,result() 以 Option 形式返回 SniffedBehavior。clear() 丢弃这些字节;协议核心在下一个字节事件开始时调用它,held 缓冲区规则保证这样做是安全的。

Passthrough<K> 是只有一个流式出站的协议的中继尾段。它跟踪两个流结束标志,并推入相应的半关闭:

方法 推入的 effect
on_transport(data, fx) Forward { key, range: 0..data.len() };返回 data.len()
on_outbound(data, fx) 无;原样 staging data,或以 staging_full() 失败
on_outbound_eof(fx) ShutdownTransport,若传输层已结束则再推入 Finish
on_transport_eof(fx) Shutdown { key },若出站已结束则再推入 Finish
on_outbound_gone(fx) ShutdownTransport、Finish(在 ConnectFailed 或 OutboundError 之后)

PassthroughCore:最小的真实协议核心

Section titled “PassthroughCore:最小的真实协议核心”

PassthroughCore<T> 中继一个在第一个字节之前就已知目的地的流。TUN 入站使用它,因为 IP 协议栈已经告知客户端要去哪里。它声明 Key = Single、Target = Flow<T>、TransportAddr = ()、STAGING_RESERVE = 0 和 BUF_SIZE = 8 * 1024。

protocols/src/core/mod.rs
impl<T> PassthroughCore<T> {
pub const BUF_SIZE: usize = 8 * 1024;
pub fn new(flow: Flow<T>) -> Self;
pub fn sniffing(flow: Flow<T>) -> Self;
pub fn is_established(&self) -> bool;
}

以 passthrough_core_opens_on_the_first_bytes_and_relays_verbatim 和 a_sniffing_passthrough_core_holds_the_prefix_and_opens_with_the_host 为线索追踪:

  1. 首批字节,不嗅探。 Transport(b"hello") 立即打开流,因为客户端发出首批字节之前运行时不会投递任何事件。effect 依次为 Open { key: Single, target }、来自 enter(Phase::Relay) 的 SetDeadline(Some(RELAY_IDLE_TIMEOUT))、来自 touch 的第二个 SetDeadline,以及 Forward { range: 0..5 }。调用返回 5。

  2. 下行。 Outbound { data: b"world" } 推入 SetDeadline(Some(RELAY_IDLE_TIMEOUT)),并原样 staging world。

  3. 两半都关闭。 TransportEof 推入 Shutdown { key: Single };随后 OutboundEof 推入 ShutdownTransport 和 Finish。

  4. 改为嗅探。 用 PassthroughCore::sniffing 和一个 IP 目的地(worth_sniffing)构建时,第一个 Transport(b"GET / HTTP/1.1\r\nHo") 只推入 SetDeadline(Some(SNIFF_TIMEOUT)),并把所有字节消费进 SniffPrefix。

  5. 找到 Host。 下一个 Transport(b"st: example.com\r\n\r\nbody") 补全了 Host 头部。协议核心推入带 sniffed.domain == "example.com" 的 Open、ForwardHeld { range: 0..41 }(全部已收集的字节,包括 body)以及 SetDeadline(Some(RELAY_IDLE_TIMEOUT))。

  6. 释放 held 缓冲区。 在运行时下,接下来的 Transport(b"more") 只会在 held 转发被应用之后到达。协议核心重新设置空闲截止时间(SetDeadline(Some(RELAY_IDLE_TIMEOUT))),清空前缀,并从事件切片转发 0..4。harness 不应用任何 effect,所以测试只检查之后 held() 为空。

如果嗅探截止时间先到,Expired::Sniff 会以 sniffed: None 打开流,并用 ForwardHeld 转发已收集的数据(a_sniffing_passthrough_core_opens_on_the_sniff_deadline)。域名目的地则完全跳过嗅探(a_sniffing_passthrough_core_with_a_domain_target_opens_at_once)。

protocols/src/core/harness.rs 中的 CoreHarness 是一个手动驱动的运行时替身。它持有协议核心、一个 EffectList、一个 WriteBuffer<HARNESS_STAGING>(HARNESS_STAGING = 64 * 1024),对于数据报协议核心还有一个 PacketList。它不做任何 I/O。

protocols/src/core/harness.rs
pub const HARNESS_STAGING: usize = 64 * 1024;
pub struct CoreHarness<C: ProxyCoreDecode> {
pub core: C,
/* private fields */
}
impl<C: ProxyCoreDecode> CoreHarness<C> {
pub fn new(core: C) -> Self;
pub fn over_datagrams(core: C) -> Self;
pub fn event(&mut self, event: Event<'_, C>) -> Result<(usize, Vec<Effect<C>>), C::Error>;
pub fn transport(&mut self, data: &mut [u8]) -> Result<(usize, Vec<Effect<C>>), C::Error>;
pub fn feed(&mut self, data: &mut [u8]) -> Result<(usize, Vec<Effect<C>>), C::Error>;
pub fn outbound(
&mut self,
key: C::Key,
data: &mut [u8],
) -> Result<(usize, Vec<Effect<C>>), C::Error>;
pub fn staged(&mut self) -> Vec<u8>;
pub fn staged_packets(&mut self) -> Vec<(Vec<u8>, C::TransportAddr)>;
pub fn held(&self, range: std::ops::Range<usize>) -> Vec<u8>;
}
方法 作用
new / over_datagrams 为字节流协议核心构建 harness,或构建一个汇集器通过 with_packets 创建、从而可以使用 stage_to 的 harness。
event 投递任意事件一次;返回消费计数和推入的 effect。
transport 即 event(Event::Transport(data))。
feed 像运行时那样投递传输层字节:只要协议核心有进展,就在未消费的尾部上再次调用。结果中 Forward 和 SendTo 的区间会被换算为相对 data 的绝对区间;held 区间保持不变。
outbound 即 event(Event::Outbound { key, data })。
staged 取走到目前为止发往传输层的全部 staging 数据。staging 数据会跨调用累积,直到被取走。
staged_packets 按顺序取走已 staging 的数据包及其对端。
held 协议核心在 range 下的 held 字节,与转发时的解析结果一致。

一个典型的单元测试,摘自 protocols/tests/unit/core/mod.rs:

protocols/tests/unit/core/mod.rs
let mut h = CoreHarness::new(PassthroughCore::sniffing(flow()));
let mut opaque = [0x16u8, 0x03, 0x01];
h.transport(&mut opaque).unwrap();
let (_, fx) = h.event(Event::Deadline).unwrap();
assert!(matches!(&fx[0], Effect::Open { target, .. } if target.sniffed.is_none()));
assert!(matches!(&fx[1], Effect::ForwardHeld { range, .. } if *range == (0..3)));
assert!(h.core.is_established());

protocols/tests/unit/core/mod.rs 还用了一种更小的模式:一个 with_sink 函数,它在 WriteBuffer::<256> 上构建 Effects::new,运行一个闭包,并返回推入的 effect 和 staging 的字节。它用来测试那些不是完整协议核心的辅助组件(Timing、Passthrough),两个不嗅探的 PassthroughCore 测试也直接通过它调用 handle。

协议核心通过从 handle 返回 Err 使连接失败。出站失败对运行时来说并不致命:它们以事件形式到达协议核心,由协议核心决定如何处理。其他所有导致连接提前结束的情况都是 RuntimeError:

concepts/src/runtime.rs
pub enum RuntimeError<E> {
Transport(io::Error),
Core(E),
UnknownKey,
DuplicateKey,
RangeOutOfBounds,
BadConsume,
FrameTooLarge,
WrongLinkKind,
StagedWithoutPeer,
}
变体 原因 Display 文本
Transport 传输层读取失败(包括数据报链路报告了一次流式读取),或字节流传输层的写、flush 或关闭失败或写入了零字节。被数据报传输层拒绝的数据包是 Event::TransportSendFailed,不是这个错误。 transport: …
Core handle 返回了 Err proxy core: …
UnknownKey 针对没有槽位的 key 的转发、发送或关闭 effect targets an unknown outbound key
DuplicateKey 在存活的 key 上 Open open reuses a live outbound key
RangeOutOfBounds 区间超出事件切片或 held 缓冲区 forward range outside the event slice or held buffer
BadConsume 消费量超过所给数据,或少于一个完整的出站负载或传输层数据报 core consumed an impossible byte count
FrameTooLarge 未解析区域占满 BUF_SIZE,而协议核心还需要更多数据 protocol frame exceeds the read buffer
WrongLinkKind 在数据报出站上执行流式 effect,或者反过来 stream effect on a datagram outbound or vice versa
StagedWithoutPeer 在数据包之外向数据报传输层 staging 了字节 bytes staged toward a datagram transport without a peer

UnknownKey、DuplicateKey、RangeOutOfBounds、BadConsume、WrongLinkKind 和 StagedWithoutPeer 都是协议核心的 bug,而非对端行为:协议核心应把畸形的线上输入转换为自己的 Error,绝不能转换为越界的 effect。FrameTooLarge 则意味着要么对端发送了比缓冲区更长的帧,要么 BUF_SIZE 对该协议来说太小。

取消由运行时负责:丢弃 ProxyServerRuntime 的 future 或 stream,就会随之关闭传输层和所有出站。协议核心对 I/O 没有任何 Drop 义务。

常量 值 位置
ProxyCoreDecode::MAX_DATAGRAM 默认值 4096 concepts/src/core.rs
INLINE_EFFECTS 4 concepts/src/core.rs
HARNESS_STAGING 64 * 1024 protocols/src/core/harness.rs
HANDSHAKE_TIMEOUT 10 s protocols/src/core/mod.rs
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs
SNIFF_LIMIT 4 * 1024 protocols/src/sniff/mod.rs

树内协议核心及其声明的常量:

协议核心 Key TransportAddr STAGING_RESERVE MAX_DATAGRAM BUF_SIZE
PassthroughCore(protocols/src/core/mod.rs) Single () 0 默认 8 * 1024
HttpCore(protocols/src/http/core.rs) Single () 256 默认 MAX_HEAD(64 * 1024)
ShadowsocksCore(protocols/src/ss_legacy/core.rs) Single () 32 + 2 * CHUNK_OVERHEAD + 28 默认 20 * 1024
Ss2022Core(protocols/src/ss_2022/core.rs) Single () 32 + 1 + 8 + 32 + 2 + 2 * TAG_SIZE + RECORD_OVERHEAD + 115 默认 32 * 1024
TrojanCore(protocols/src/trojan/core.rs) FlowKey () PACKET_HEADER_MAX.next_multiple_of(16) + downlink_overhead(Self::BUF_SIZE) MAX_LENGTH(8192) 16 * 1024
VlessCore(protocols/src/vless/core.rs) FlowKey () 272 + downlink_overhead(Self::BUF_SIZE) 8192 16 * 1024
VMessCore(protocols/src/vmess/core.rs) FlowKey () 4096 8192 32 * 1024
Hy2StreamCore(protocols/src/hysteria/server/inbound.rs) Single () 2048 默认 8 * 1024
Hy2UdpCore(protocols/src/hysteria/server/datagrams.rs) u32 () 4096 MAX_UDP_SIZE(4096) 16 * 1024
TunUdpCore(protocols/src/tun/udp.rs) Single SocketAddr 4096 4096 8 * 1024

使用 FlowKey 的协议核心承载 mux.cool 子流;protocols/src/mux/demux.rs 中的 downlink_overhead 计入了一次出站读取可能被拆分成的多个 Keep 帧的头部。

测试 文件 固定的行为
core_is_driven_without_any_io concepts/tests/runtime.rs 协议核心可以用手工构造的事件运行;被截断的帧保持未消费;原地解密
relays_two_keys_and_completes_on_fin concepts/tests/runtime.rs 排在 Open 之后的转发等待连接完成;Finish 以 Traffic 总量完成
stalled_outbound_holds_uplink_but_not_other_downlink concepts/tests/runtime.rs 被阻塞的转发挡住其后的一切;其他下行照常流动
connect_failure_reaches_the_core_as_an_event concepts/tests/runtime.rs ConnectFailed 的投递
half_close_propagates_both_ways_and_finishes concepts/tests/runtime.rs Shutdown、OutboundEof、ShutdownTransport
deadline_event_lets_the_core_time_out、deadline_is_armed_against_the_tokio_clock concepts/tests/runtime.rs 单一截止时间与 tokio 时钟
frame_larger_than_the_buffer_is_an_error concepts/tests/runtime.rs FrameTooLarge
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites、a_held_range_past_the_buffer_is_rejected concepts/tests/runtime.rs held 缓冲区规则、RangeOutOfBounds
datagram_outbound_is_never_truncated_by_staging_backpressure concepts/tests/runtime.rs 读取数据报前须有容纳整个数据包的空间
datagram_transport_demultiplexes_peers_and_frames_replies、staging_without_a_peer_under_a_datagram_transport_is_an_error、stage_to_is_refused_under_a_stream_transport concepts/tests/runtime.rs 数据报传输层与 StagedWithoutPeer
a_refused_datagram_send_keeps_the_key_alive、a_refused_transport_packet_is_reported_not_fatal、single_peer_datagram_transport_demuxes_sessions_and_reports_refusals concepts/tests/runtime.rs SendFailed 和 TransportSendFailed 不致命
stream_mode_reports_each_unit_of_work concepts/tests/runtime.rs 每一步的 Traffic 增量;非负载步骤的增量为零
timing_arms_handshake_once_then_idle_per_byte_event、timing_reports_handshake_and_sniff_expiry_to_the_core protocols/tests/unit/core/mod.rs Timing
sniff_prefix_finds_a_host_across_pushes_and_keeps_the_bytes、sniff_prefix_takes_no_more_than_its_budget protocols/tests/unit/core/mod.rs SniffPrefix 与 SNIFF_LIMIT
passthrough_half_closes_each_side_and_finishes_on_the_second protocols/tests/unit/core/mod.rs Passthrough
passthrough_core_*(2 个测试,通过 with_sink)、a_sniffing_passthrough_core_*(3 个测试,通过 CoreHarness) protocols/tests/unit/core/mod.rs PassthroughCore:打开、中继、半关闭、失败、空闲过期和嗅探
stream_sessions_open_forward_and_end_with_fresh_generations protocols/tests/unit/mux/demux.rs 带 generation 标记的 SubKey
a_frame_straddling_chunks_is_held_and_forwarded_from_the_held_buffer、vmess_keeps_every_frame_one_read_completes protocols/tests/unit/mux/demux.rs 跨 VMess 分块拆分的 mux 帧从 held 缓冲区转发
a_chunk_split_across_reads_decodes_its_header_once protocols/tests/unit/vmess/core.rs 绝不重复解密头部(VMess)
chunk_encoder_and_decoder_agree_and_never_decrypt_twice protocols/tests/unit/ss_legacy/aead.rs 绝不重复解密头部(Shadowsocks ChunkDecoder)

用 cargo test -p etemenanki-concepts --test runtime 和 cargo test -p etemenanki-protocols --lib core::tests 运行它们(该过滤条件也会选中每个协议核心自己的 core::tests 模块)。