跳转到内容

客户端 codec 与客户端运行时

源码文件:31 个 · 核对版本 Etemenanki 596916d · katana v3.0.1
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/http/codec.rs
  • Etemenanki/protocols/src/socks/codec.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/trojan/codec.rs
  • Etemenanki/protocols/src/vless/codec.rs
  • Etemenanki/protocols/src/vmess/codec.rs
  • Etemenanki/protocols/src/ss_legacy/codec.rs
  • Etemenanki/protocols/src/ss_2022/codec.rs
  • Etemenanki/protocols/tests/support/pipeline.rs
  • Etemenanki/protocols/tests/pipeline/http.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/protocols/tests/pipeline/trojan.rs
  • Etemenanki/protocols/tests/pipeline/vless.rs
  • Etemenanki/protocols/tests/pipeline/vmess.rs
  • Etemenanki/protocols/tests/pipeline/shadowsocks.rs
  • Etemenanki/protocols/tests/unit/socks/codec.rs
  • Etemenanki/protocols/tests/unit/socks/protocol.rs
  • Etemenanki/protocols/tests/unit/vless/codec.rs
  • Etemenanki/protocols/tests/unit/vmess/codec.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • katana/src/outbound/mod.rs
  • katana/src/outbound/proxy.rs

服务端协议核心解析客户端发来的内容,并决定每条流去往哪里。当一条流要交给另一台代理服务器、而不是直接去往目的地时,出站路径上就需要有东西来说那台代理的协议:写出请求头,协议有应答时等待应答,把每个载荷字节封装成线上帧,再打开传回来的帧。在 Etemenanki 中,这就是模型的客户端一侧。

客户端一侧分三层,都在 etemenanki-concepts 中:

  • codec:针对单条流的 sans-I/O 状态机,每个协议在 etemenanki-protocols 中各写一份;
  • ProxyClientRuntime:在一条已拨通的上游连接上驱动一个 codec,从外部看就是一个普通字节流(或一个数据报 link);
  • ProxyClientConnector:一个 Connector,为每条流构造一个 codec 和一个运行时,使服务端运行时可以经由上游代理拨号。

本页面向修改这三层中任意一层、或新增客户端协议的贡献者,并假定你已经从服务端协议核心和服务端运行时了解了服务端一侧。

层 负责 不负责
codec(ProxyCoreEncodeHandshake 加上 ProxyCoreEncode 或 ProxyCoreEncodeDatagram) 协议本身:请求头、握手应答、分帧、加密、关闭。 socket、缓冲区、waker、定时器。它只能看到字节切片和一个 Staging 区域。
ProxyClientRuntime(concepts/src/client.rs) 上游 stream、发往上游的 staging 缓冲区、从上游读取的读缓冲区、握手循环、明文一侧的 AsyncRead/AsyncWrite 或 DatagramLink,以及检查 codec 是否遵守约定。 拨号(它调用内层 Connector)、选择 codec、决定“已连接”对服务端协议核心意味着什么。
ProxyClientConnector / ProxyClientConnecting(concepts/src/client.rs) 按流选择 stream codec 或数据报 codec,并且只在上游已拨通、握手已完成、握手字节已 flush 之后才让拨号 resolve。 拨号 resolve 之后的一切:从那时起,服务端运行时直接 poll ProxyClientRuntime。

客户端一侧不 spawn 任何 task,也不使用 channel。ProxyClientRuntime 由持有它的一方 poll,通常是把它存为某个出站的 ProxyServerRuntime,因此 codec、两个缓冲区和上游 stream 都在服务端运行时的 task 内被 poll。唯一的例外位于客户端运行时之下:gRPC 客户端传输层把它的 HTTP/2 连接驱动作为单独 spawn 的 task 运行,stream 被 drop 时该 task 会被 abort。

每个客户端 codec 都实现握手 trait。载荷形态是叠加在它之上的第二个 trait,因此 codec 恰好只实现其流所具有的操作。

concepts/src/core.rs
pub trait ProxyCoreEncodeHandshake {
type Target;
type Error;
const STAGING_RESERVE: usize;
fn start(&mut self, out: &mut Staging<'_>) -> Result<Handshake, Self::Error>;
fn reply(&mut self, wire: &mut [u8], out: &mut Staging<'_>) -> Result<Reply, Self::Error>;
fn finish(&mut self, out: &mut Staging<'_>) -> Result<(), Self::Error>;
}
项 含义
Target 运行时的内层 Connector 拨号的对象:上游代理服务器,而不是流的目的地。目的地保存在 codec 内部,由 start 编码。etemenanki-protocols 中的所有 codec 都使用 Destination。
Error codec 的错误类型。运行时要求 Error + Send + Sync + 'static,并把它包装成 kind 为 InvalidData 的 io::Error。
STAGING_RESERVE 一次 start、reply、finish 或 seal 调用在所接收明文之外最多写入 staging 的字节数:头部、长度前缀、tag、填充。运行时给流 seal 至少 STAGING_RESERVE + 1 字节的空间,给报文 seal STAGING_RESERVE + plain.len() 字节,因此遵守这一上限的 codec 永远不会空间不足。
start 拨号成功后立即调用一次。它写入请求头,并说明在明文开始流动之前是否必须先解析一个应答。
reply 仅在握手处于 AwaitReply 时调用,参数是尚未解析的上游字节(可变,以便 codec 原地解密)。它可以写入下一轮的内容。最后一个应答之后剩余的字节就是最初的线上帧。
finish 明文一侧已半关闭。codec 写入用来告知上游的内容:一个关闭帧;如果线路自身的 EOF 就能表达,则什么都不写。

握手状态由两个小枚举承载:

concepts/src/core.rs
pub enum Handshake {
Done,
AwaitReply,
}
pub enum Reply {
NeedMore,
Step { consumed: usize, next: Handshake },
}

start 返回 Handshake::Done 表示协议发出头部后就直接开始(VMess、VLESS、Trojan、Shadowsocks)。Handshake::AwaitReply 表示上游必须先应答(SOCKS 方法选择、HTTP CONNECT 状态行)。Reply::NeedMore 要求更多线上字节。Reply::Step 表示前 consumed 个字节包含一个完整应答,next 说明是否还有下一轮。

Step 必须取得进展:消费字节、写入字节,或者结束握手。三者都不做的 step 会被运行时拒绝(见握手循环)。

concepts/src/core.rs
pub trait ProxyCoreEncode: ProxyCoreEncodeHandshake {
fn seal(&mut self, plain: &[u8], out: &mut Staging<'_>) -> Result<usize, Self::Error>;
fn open(&mut self, wire: &mut [u8]) -> Result<Opened, Self::Error>;
}
pub enum Opened {
Frame { consumed: usize, plain: Range<usize> },
NeedMore,
End { consumed: usize },
}
  • seal 取 plain 的一个前缀,把它写成线上帧,返回取走的明文字节数。这个数可以少于提供的量,例如只取一帧的量,但至少为 1。
  • open 查看尚未解析的线上字节,并原地打开最前面的一帧。Frame { consumed, plain } 表示前 consumed 个字节已处理完毕,而 plain 是其中的一个子区间,现在存放着解密后的明文。对控制帧来说 plain 为空。End { consumed } 表示上游用一个关闭帧干净地结束了这条流。NeedMore 要求更多字节。Frame 或 End 必须至少消费 1 个字节,且不超过整个切片。

明文为空的帧是 codec 在不增加应答轮次的情况下消费响应头的方式。VLESS 响应头、VMess 响应头、旧版 Shadowsocks 的响应 salt,以及 Shadowsocks 2022 的响应 salt 和固定响应头,都会作为空帧打开,由运行时跳过。

数据报 codec 在通往上游的一条流上传输带有逐包对端地址的报文(Trojan、VLESS 和 VMess 的 UDP)。

concepts/src/core.rs
pub type OpenedFrom = (Opened, Option<Destination>);
pub trait ProxyCoreEncodeDatagram: ProxyCoreEncodeHandshake {
fn seal_to(
&mut self,
plain: &[u8],
to: &Destination,
out: &mut Staging<'_>,
) -> Result<Option<()>, Self::Error>;
fn open_from(&mut self, wire: &mut [u8]) -> Result<OpenedFrom, Self::Error>;
}

一次 seal_to 就是一个报文,要么整个取走,要么在放不进所给空间时返回 Ok(None) 且不写入任何内容。open_from 打开的一个数据帧就是一个报文,连同其来源一起返回。来源为 None 的帧不是报文(VLESS 或 VMess 的响应头,或一个空的 VMess 分块),运行时会消费并跳过它。open_from 遵循与 open 相同的进展规则。

对端地址的用法取决于协议。TrojanDatagram 把 to 写进每个报文,并从每个报文中读出来源。VlessDatagram 和 VMessDatagram 在请求头中固定目标,忽略 to,并把该目标报告为它们打开的每个报文的来源。

concepts/src/core.rs
pub enum NoCodec<Target, Error> {
Never(Infallible, PhantomData<(Target, Error)>),
}

NoCodec 是不可构造的类型,以 STAGING_RESERVE = 0 实现了全部三个 trait。ProxyClientConnector 总是同时带有一个 stream codec 类型和一个数据报 codec 类型,所以只承载其中一种的协议会把另一种指定为 NoCodec。app 和 katana 都把它别名为 NoUdp = NoCodec<Destination, io::Error>,并用于 HTTP、SOCKS CONNECT 和两种 Shadowsocks。

Codec 文件 start 返回 finish 写入 STAGING_RESERVE
HttpConnect protocols/src/http/codec.rs AwaitReply(一轮:状态 200,其他都是错误) 无 HttpConnect::REQUEST_MAX(1024)
SocksConnect protocols/src/socks/codec.rs AwaitReply(方法、可选的用户名/密码、请求) 无 528
TrojanStream、TrojanDatagram protocols/src/trojan/codec.rs Done 无 RESERVE(REQUEST_HEADER_MAX)
VlessStream、VlessDatagram protocols/src/vless/codec.rs Done 无 RESERVE(REQUEST_HEADER_MAX)
VMessStream、VMessDatagram protocols/src/vmess/codec.rs Done 一个终止分块 HEADER_MAX.next_multiple_of(64)
SsStream protocols/src/ss_legacy/codec.rs Done 无 32 + CHUNK_OVERHEAD + AddressCodec::MAX_LEN + 61
Ss2022Stream protocols/src/ss_2022/codec.rs Done 无 2048

start 返回 Done 的 codec 对 reply 一律返回错误,因为运行时从不对它们调用 reply。线上格式见各协议的页面。

concepts/src/client.rs
pub struct ProxyClientRuntime<const BUF_SIZE: usize, Codec, Conn>
where
Codec: ProxyCoreEncodeHandshake,
Conn: Connector<Codec::Target>,
{ /* 私有字段 */ }
impl<const BUF_SIZE: usize, Codec, Conn> ProxyClientRuntime<BUF_SIZE, Codec, Conn>
where
Codec: ProxyCoreEncodeHandshake,
Codec::Error: Error + Send + Sync + 'static,
Conn: Connector<Codec::Target>,
{
pub fn new(codec: Codec, connector: &mut Conn, target: Codec::Target) -> Self;
pub fn codec(&self) -> &Codec;
}

new 断言 BUF_SIZE > Codec::STAGING_RESERVE,否则 panic,信息为 BUF_SIZE must exceed the codec's STAGING_RESERVE or nothing can be sealed。随后它调用 connector.connect(target),并把返回的 future 固定在一个 Box 中。此时还不会 poll 任何东西。拨号 future 在任一明文侧方法第一次被 poll 时才开始运行,或者在包装该运行时的 ProxyClientConnecting 第一次被 poll 时运行。

运行时实现哪些 trait 取决于 codec:

Codec 约束 实现的 trait 用作
Codec: ProxyCoreEncode AsyncRead + AsyncWrite 流出站
Codec: ProxyCoreEncodeDatagram DatagramLink<Addr = Destination> 数据报出站
始终 Unpin DatagramLink: Unpin 以及 Connector::Stream: AsyncRead + AsyncWrite + Unpin 约束的要求

Unpin 是无条件实现的。这是 sound 的,因为拨号 future 以 Pin<Box<_>> 形式保存,上游 stream 由 Connector::Stream 约束保证为 Unpin,而 codec 从不被 pin。

私有状态:

字段 类型 含义
codec Codec 协议状态机。
wire Wire<Conn::Future, Conn::Stream> Connecting(Pin<Box<F>>)、Up(S) 或 Down。Down 是终态:之后每个需要线路的调用都以 NotConnected 失败(空的 poll_write 仍返回 Ok(0))。
staging WriteBuffer<BUF_SIZE> 已封装但尚未写入上游的字节。
down ReadBuffer<BUF_SIZE> 从上游读到的字节。未解析区域是 codec 尚未打开的部分。
opened Option<(usize, Range<usize>)> 已打开但尚未全部交给读取方的帧:它的线上长度,以及仍待交付的明文,均为 down 中的绝对偏移。
wire_read_closed bool 上游读取一侧已返回 EOF。
ended bool codec 报告了 Opened::End。
finished bool finish 已写入 staging。它只写入一次。
handshaking bool start 或上一次 reply 要求的应答尚未解析。

ProxyClientConnector 与 ProxyClientConnecting

Section titled “ProxyClientConnector 与 ProxyClientConnecting”
concepts/src/client.rs
pub struct ProxyClientConnector<const BUF_SIZE: usize, Make, Conn, Upstream> {
make: Make,
dialer: Conn,
upstream: Upstream,
}
impl<const BUF_SIZE: usize, Make, Conn, Upstream>
ProxyClientConnector<BUF_SIZE, Make, Conn, Upstream>
{
pub fn new(make: Make, dialer: Conn, upstream: Upstream) -> Self;
}
impl<const BUF_SIZE: usize, Make, Conn, Upstream, Flow, S, D> Connector<Flow>
for ProxyClientConnector<BUF_SIZE, Make, Conn, Upstream>
where
Make: FnMut(Flow) -> Outbound<S, D>,
S: ProxyCoreEncode<Target = Upstream>,
S::Error: Error + Send + Sync + 'static,
D: ProxyCoreEncodeDatagram<Target = Upstream>,
D::Error: Error + Send + Sync + 'static,
Conn: Connector<Upstream>,
Upstream: Clone,
{
type Stream = ProxyClientRuntime<BUF_SIZE, S, Conn>;
type Datagram = ProxyClientRuntime<BUF_SIZE, D, Conn>;
type Future = ProxyClientConnecting<BUF_SIZE, S, D, Conn>;
fn connect(&mut self, flow: Flow) -> Self::Future;
}
pub struct ProxyClientConnecting<const BUF_SIZE: usize, S, D, Conn>
where
S: ProxyCoreEncodeHandshake,
D: ProxyCoreEncodeHandshake<Target = S::Target>,
Conn: Connector<S::Target>,
{
runtime: Option<
Outbound<ProxyClientRuntime<BUF_SIZE, S, Conn>, ProxyClientRuntime<BUF_SIZE, D, Conn>>,
>,
}

Flow 是服务端一侧用来路由的类型:在 app 中是 Flow,在 katana 中是 Destination,在 concepts 测试中是 String。Upstream 是上游代理的地址,每次拨号都会克隆一份。ProxyClientConnecting 实现了 Future,其 Output = io::Result<Outbound<ProxyClientRuntime<BUF_SIZE, S, Conn>, ProxyClientRuntime<BUF_SIZE, D, Conn>>>。

一个客户端运行时恰好分配两个装箱的 [u8; BUF_SIZE] 数组(通过 boxed_array,因此即使 BUF_SIZE 很大也不会经过栈)和一个装箱的拨号 future。缓冲区从不增长,所以背压来自缓冲区被填满。

flowchart LR
  W["poll_write(plain)"] --> S["codec.seal"]
  S --> ST["staging: WriteBuffer"]
  ST --> U["上游 stream"]
  U --> D["down: ReadBuffer"]
  D --> O["codec.open,原地解密"]
  O --> R["poll_read 复制出明文"]
  • staging 是先进先出的。start、reply、seal、seal_to 和 finish 都通过 WriteBuffer::staging() 追加。drain 把 pending() 写入上游并推进 start。尾部满了时,WriteBuffer::room() 把待发送字节滑到最前面,这是该缓冲区自行进行的唯一一次复制。因此 start 写入的请求头总是先于第一个封装帧发出。
  • down 有三个区域:已打开的字节、未解析的 [start, end),以及空闲空间。fill 读入空闲尾部,并且只在尾部为空时调用 compact(),把未解析字节移到最前面。如果未解析区域占满整个数组(is_saturated()),说明这一帧放不下,fill 失败。
  • 原地处理。open 和 open_from 在 down 内解密。运行时直接把明文从 down 复制到调用方的 ReadBuf,没有中间缓冲区。通过较小的 ReadBuf 读取大帧时,会分多次调用交付:opened 记住仍待交付的绝对区间,而 down.advance_start(consumed) 只在最后一个字节交付之后才执行。
stateDiagram-v2
  [*] --> Connecting: new
  Connecting --> Down: 拨号错误或得到数据报 link
  Connecting --> Down: start 出错
  Connecting --> Handshaking: start 返回 AwaitReply
  Connecting --> Open: start 返回 Done
  Handshaking --> Handshaking: Step,next 为 AwaitReply
  Handshaking --> Open: Step,next 为 Done
  Handshaking --> Down: EOF 或线路 I/O 错误
  Open --> Ended: codec 打开 End
  Open --> Down: 线路 I/O 错误
  Ended --> [*]
  Down --> [*]

Connecting、Up 和 Down 是 Wire 的变体。Handshaking 是 handshaking = true 的 Up,Open 是握手已完成的 Up,Ended 是 ended = true 的 Open。is_ready() 恰好在 Open 和 Ended 时为 true。在此之前,只要拨号或握手仍在进行,所有明文操作都返回 Pending。

ready() 在任何明文侧 poll 中驱动这台状态机:wire() 完成拨号并运行 start,然后 handshake() 运行各轮应答。因为 poll_write 和 poll_read 都会调用它,只写不读的调用方也能让握手应答被读取。concepts/tests/client.rs 中的 plaintext_waits_for_every_handshake_reply 确保只写的调用方不会死锁。

handshake() 在 handshaking 为 true 时循环:

  1. drain 所有已写入 staging 的内容(start 的问候,或上一次 reply 写入的那一轮)。在上游接收全部内容之前,它返回 Pending。
  2. 如果 down 中有未解析字节,调用 codec.reply(down.unparsed(), staging):
    • Step { consumed, next } 的 consumed 大于给出的字节数时,以 codec consumed more reply than it was given 失败。
    • Step 的 consumed == 0、没有写入任何内容(通过比较调用前后的 staging.room() 判断)且 next == AwaitReply 时,以 codec handshake step made no progress 失败。没有这项检查,循环会空转。
    • 否则运行时调用 down.advance_start(consumed),设置 handshaking = (next == AwaitReply) 并继续循环。
    • NeedMore 进入第 3 步。
  3. 从上游 fill 一块数据。此处遇到 EOF 会让线路进入 Down,并以 UnexpectedEof、upstream closed during the handshake 失败。

与最后一个应答在同一次读取中到达的字节留在 down 中,作为最初的帧被打开。plaintext_waits_for_every_handshake_reply 在一次写入中发送连接应答和一个数据帧,并把该帧读回。

poll_write(plain):

  1. 空的 plain 立即返回 Ok(0),不驱动拨号。
  2. ready():拨号、start、各轮应答。
  3. make_room(STAGING_RESERVE + 1):把已写入 staging 的字节写往上游,直到空出这么多空间。这里的 Pending 就是背压路径:上游不接收字节,调用方的写入就保持 pending。
  4. 向 seal 提供 min(plain.len(), room - STAGING_RESERVE) 个字节。返回 0 或超过所提供的数量时,以 codec sealed an impossible byte count 失败。
  5. 尝试 drain。Pending 的 drain 会被忽略,因为字节已经写入 staging,下一次 poll 或 poll_flush 会把它们发出。错误则会返回。
  6. 返回 Ok(taken)。

poll_flush 运行 ready(),drain 所有已写入 staging 的内容,然后 flush 上游 stream。poll_shutdown 写入一次 finish(在 make_room(STAGING_RESERVE) 之后),然后 drain、flush,最后关闭上游 stream 的写入一侧。handshake_then_seal_and_open_over_the_wire 检查明文侧 shutdown 会把关闭帧放上线路,然后再半关闭线路。

服务端运行时在每次取走某个 forward 至少 1 个字节的 poll_write 之后都会跟一次 poll_flush。如果 flush 返回 Pending,它会把该 slot 标记为 need_flush,并在下一次处理该 key 时再次 flush,因此封装好的帧不会滞留在 staging 中。

poll_read(buf) 首先运行与 flush 相同的前置步骤(ready() 然后 drain),区别只有一点:Pending 的 drain 不会阻塞读取。代码中的注释是这样说的:“a stalled wire write must not block reading”。如果握手尚未完成,poll_read 返回 Pending。否则进入循环:

  1. 如果有一个帧处于 opened,按 buf 能容纳的量复制其明文并返回。整个帧交付完毕后,清除 opened,并把 down 推进到该帧之后。
  2. 如果 ended,返回 Ok(()) 且不填充任何内容,即 EOF。
  3. 如果 down 中有未解析字节,调用 codec.open:
    • Frame 需通过 check_frame。plain 为空时跳过(消费该帧并继续循环)。否则由 opened 记录该帧。
    • End 以空区间通过 check_frame,被消费,并设置 ended。
    • NeedMore 进入下一步。
  4. fill。遇到上游 EOF 时,若 down 为空则返回 EOF;若还剩部分帧,则返回 UnexpectedEof、upstream closed inside a frame。

check_frame(len, consumed, plain) 拒绝 consumed == 0、consumed > len、plain.start > plain.end 和 plain.end > consumed,报错 codec opened an impossible frame。这些情况中的任何一种都会让循环空转或越出帧读取。

流如何结束:

上游的行为 poll_read 返回
发送一个 codec 打开为 End 的帧 EOF。之后的读取也返回 EOF,关闭帧之后的内容都不再打开。
在两帧之间关闭连接 EOF。对于 finish 不写入任何内容的 codec,这就是流的正常结束。
在一帧中间关闭连接 UnexpectedEof、upstream closed inside a frame
在握手期间关闭连接 UnexpectedEof、upstream closed during the handshake,线路进入 Down

使用 ProxyCoreEncodeDatagram codec 时,运行时是一个以 Destination 寻址的 DatagramLink。报文仍然在那一条上游 stream 上传输。

  • poll_send_to(plain, to) 运行 ready(),然后 make_room(STAGING_RESERVE + plain.len())。如果 staging 已经为空而空间仍然不够,这个报文永远放不下,调用以 InvalidInput、frame larger than the client runtime's buffer 失败。如果空间足够而 codec 仍返回 None,调用以 InvalidInput、codec refused a packet that fits its reserve 失败。成功时它会尝试 drain,并返回 Ok(plain.len()),因为报文要么整个被取走,要么完全不取。
  • poll_recv_from(buf) 每次调用返回一个报文。来源为 Some 的数据帧被复制到 buf 中,buf 较小时像 UDP 一样截断。无论哪种情况,整个帧都会被消费。来源为 None 的帧被跳过。在 Opened::End 之后,或在上游 EOF 且 down 为空时,调用以 UnexpectedEof 失败。数据报 link 没有“流结束”的概念,所以服务端协议核心会把结束视为该 key 的出站错误。

在服务端运行时中,发送失败变成 Event::SendFailed,key 保持存活;接收失败则变成 Event::OutboundError,并移除该 key。

connect(flow) 同步完成三件事:

  1. 克隆 upstream。

  2. 调用 make(flow)。闭包返回 Outbound::Stream(codec) 或 Outbound::Datagram(codec),由此决定这条流的载荷形态。

  3. 用 ProxyClientRuntime::new(codec, &mut dialer, upstream) 包装 codec(这一步创建拨号 future),并把它放在 ProxyClientConnecting 中返回。

ProxyClientConnecting 是一个 future。每次 poll 都会调用运行时的私有方法 poll_connected:先 ready()(拨号、start、每一轮应答),再 drain()(写出所有握手字节),然后对上游 stream 调用 poll_flush。只有三者都完成,future 才 resolve 为 Ok(Outbound::Stream(runtime)) 或 Ok(Outbound::Datagram(runtime))。出错时它会 drop 运行时,上游连接也随之被 drop,然后 resolve 为该错误。flush 出错还会在运行时被 drop 之前让线路进入 Down。future resolve 之后再次 poll 会 panic,信息为 polled after completion。

正是这个时序让服务端协议核心收到的事件名副其实。服务端运行时把 connector future 的结果转换成 Event::Connected 或 Event::ConnectFailed(见服务端运行时)。因此:

  • Connected 只在上游已拨通、codec 的握手已完成、握手字节已写入并 flush 到上游 stream 之后才送达。在收到 Connected 时回复自己客户端的 SOCKS 或 HTTP 服务端协议核心,因此只会在上游代理同意之后才回复。
  • ConnectFailed 涵盖此前的所有失败:拨号本身失败、SOCKS 方法或请求被拒、HTTP 状态不是 200,或上游在握手期间关闭。协议核心可以向客户端发送真正的拒绝,而不是过早的成功。

在 slot 仍处于连接状态时,drive_effects 会在服务端运行时 effect 队列中第一个指向该 key 的 forward、send 或 shutdown 处停止应用(LinkState::Connecting(_) => break);排在它后面的 effect 也一并等待。因此,协议核心与 Open 一起推入的 forward 会等到连接 future resolve 之后才执行,而请求头总是在第一个载荷帧被封装之前单独 flush 出去。

两个程序都把每个代理出站包装在 ProxyClient<const BUF: usize, S, D> 中(app/src/outbound/proxy.rs 和 katana 的 src/outbound/proxy.rs):

app/src/outbound/proxy.rs
pub type NoUdp = NoCodec<Destination, io::Error>;
pub type Make<S, D> = Box<dyn FnMut(Flow) -> link::Outbound<S, D> + Send>;
pub struct ProxyClient<const BUF: usize, S, D> {
inner: Mutex<ProxyClientConnector<BUF, Make<S, D>, TransportConnector, Destination>>,
}
impl<const BUF: usize, S, D> ProxyClient<BUF, S, D>
where
S: ProxyCoreEncode<Target = Destination, Error = io::Error>,
D: ProxyCoreEncodeDatagram<Target = Destination, Error = io::Error>,
{
pub fn new(make: Make<S, D>, transport: TransportConnector, server: Destination) -> Self;
pub fn connect(&self, flow: Flow) -> ProxyClientConnecting<BUF, S, D, TransportConnector>;
}
  • 内层 connector 是 TransportConnector(protocols/src/transports/connect.rs)。它通过 TCP、TLS、WebSocket 或 gRPC 拨号到上游的 Destination。它的 Datagram 类型是 NoDatagram:它总是产出一个字节流,因为代理的 UDP 就承载在这条流里。
  • parking_lot::Mutex 只在 connect 构造 codec 和拨号 future 期间持有。future 在锁外 await。
  • Make 闭包根据流选择 codec。对 Trojan、VLESS 和 VMess,flow.destination.network == DialNetwork::Udp 时返回数据报 codec,否则返回 stream codec。HTTP、SOCKS CONNECT 和 Shadowsocks 总是返回 stream codec。对于 UDP 流,HTTP 和 Shadowsocks 出站会在 connect_datagram 中直接返回错误,不进行拨号。SOCKS UDP 不经过客户端运行时:SocksOutbound 自己构建 UDP ASSOCIATE link(见 SOCKS UDP:SocksUdpLink)。
  • proxy_stream 和 proxy_datagram 把 resolve 得到的运行时装箱为 OutboundStream::Proxy 或 OutboundDatagram。类型不对的运行时会报错(a TCP flow was dialed as datagrams、a UDP flow was dialed as a stream)。

在 katana 中,闭包接收的是 Destination 而不是 Flow(由 dest.network == DialNetwork::Udp 选择数据报 codec),并且没有 Trojan 出站。其余完全相同。

SOCKS5 UDP 套不进 ProxyCoreEncodeDatagram codec。它的包并不承载在通往上游的流里,而是作为独立的 UDP 数据报,发往服务器在应答中给出的中继地址。因此 SocksOutbound::connect_datagram 用 TransportConnector::dial 拨号到 SOCKS 服务器,再把这条控制流交给 SocksUdpLink::associate(protocols/src/socks/udp_link.rs)。这个 link 自己实现了 DatagramLink<Addr = Destination>。

protocols/src/socks/udp_link.rs
impl<S: AsyncRead + AsyncWrite + Unpin> SocksUdpLink<S> {
pub async fn associate(
control: S,
auth: Option<(&str, &str)>,
bind: impl FnOnce(&SocketAddr) -> io::Result<UdpSocket>,
) -> io::Result<Self>;
pub fn relay(&self) -> SocketAddr;
}
  • 握手。 associate 在 control 上依次完成方法协商、可选的 RFC 1929 用户名/密码认证和 UDP ASSOCIATE 请求,之后在整个 UDP 关联期间保持 control 打开。应答状态非零时以 ConnectionRefused 失败,消息为 server rejects request: <status>。状态以十进制打印:Etemenanki 自己的 SOCKS 入站在无法把关联绑定到其客户端时应答 0x02,link 会把它报告为 server rejects request: 2。
  • 不声明来源。 请求中的 DST.ADDR 和 DST.PORT 全为零(encode_request(CMD_UDP_ASSOCIATE, None))。RFC 1928 规定客户端不知道自己的来源地址时发送全零。link 要等应答给出中继地址后才绑定 socket,而且在 NAT 之后,本地地址本来就不对。会校验来源的服务器(例如 Etemenanki 自己的 SOCKS 入站)会把关联绑定到控制连接的来源 IP,以及它中继的第一个数据报的端口。服务端的做法见 SOCKS。
  • Socket 地址族。 bind 会收到中继地址,以便据此选择地址族。app 和 katana 都绑定到中继地址族的未指定地址。
  • 应答来源。 poll_recv_from 只接受来自服务器所给中继地址的数据报。其他数据报无论格式多么正确,都会在解析之前被跳过。比较使用 protocols/src/socks/protocol.rs 中的 endpoint,也就是服务端所用的同一个辅助函数:取规范形式的 IP(IPv4 映射的 IPv6 地址按其 IPv4 地址计)和端口,忽略 IPv6 的 flow info 和 scope。因此双栈 socket 仍然能收到 IPv4 中继的应答。katana 3.0.1 基于 etemenanki-protocols 2.0.1,那个版本的 link 精确比较 socket 地址;只要 socket 按中继的地址族绑定(katana 正是这样绑定的),两种比较的结果一致。来自中继、但无法解析为中继包的数据报同样会被跳过。
  • 控制流。 每次 poll_send_to 和 poll_recv_from 都会先读空控制流,并丢弃读到的内容。控制流结束后,每次调用都以 BrokenPipe 失败,消息为 socks: the control connection closed。控制流上的读错误按原样返回,之后的每次调用都以 BrokenPipe 失败。
测试 文件 固定的行为
new_server_vs_new_client_udp protocols/tests/pipeline/socks.rs SocksUdpLink 对接真实的 SOCKS 入站,包括一个 1400 字节的包
udp_link_ignores_datagrams_not_from_the_relay protocols/tests/pipeline/socks.rs 来自其他 socket 的格式正确的应答被丢弃,而中继的应答连同其头部中的目标地址一起交付
udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay protocols/tests/pipeline/socks.rs 在绑定到 [::] 的 socket 上,IPv4 中继的应答以 IPv4 映射地址到达,仍被接受
endpoint_sees_through_ipv4_mapping_and_ignores_flow_info protocols/tests/unit/socks/protocol.rs endpoint 把 [::ffff:1.2.3.4]:5 与 1.2.3.4:5 视为相同,区分 ::1 和 127.0.0.1,并忽略 flow info

数据流:经由客户端运行时的中继

Section titled “数据流:经由客户端运行时的中继”

concepts/tests/client.rs 中的 server_runtime_relays_through_a_client_runtime_in_one_task 构建的正是这种形态。一个透传的服务端协议核心读取 4 字节目标,然后中继其余内容。它的 connector 是基于 TinyCodec 的 ProxyClientConnector。TinyCodec 是一个玩具 codec,头部为 [0xC0][len:u8][target],帧为经过 XOR 掩码的 [kind][len:u16][payload],其中 kind 1 为数据,kind 2 为关闭。服务端协议核心在 Connected 时回复 +,在 ConnectFailed 时回复 -。

sequenceDiagram
  participant C as 客户端
  participant SR as ProxyServerRuntime
  participant Core as 服务端协议核心
  participant CR as ProxyClientRuntime
  participant U as 上游代理
  Note over SR,CR: 同一个 task,没有 channel
  C->>SR: "dst1" + 载荷
  SR->>Core: Event::Transport
  Core-->>SR: Open(target), Forward
  SR->>CR: ProxyClientConnector::connect(flow) 构造 CR
  SR->>CR: poll ProxyClientConnecting
  CR->>U: codec.start 生成的头部,drain 并 flush
  CR-->>SR: 连接 future resolve 为 Ok
  SR->>Core: Event::Connected
  Core-->>C: "+"
  SR->>CR: poll_write(payload), poll_flush
  CR->>U: 封装好的 DATA 帧
  U->>CR: DATA 帧
  SR->>CR: poll_read(按 key 的 waker 被唤醒)
  CR-->>SR: 原地打开的明文
  SR->>Core: Event::Outbound
  Core-->>C: 写入 staging 的明文
  C->>SR: EOF
  SR->>Core: Event::TransportEof
  Core-->>SR: Effect::Shutdown
  SR->>CR: poll_shutdown
  CR->>U: codec.finish 生成的 CLOSE 帧,然后关闭写入一侧
  U->>CR: CLOSE 帧
  CR-->>SR: poll_read 返回 EOF(Opened::End)
  SR->>Core: Event::OutboundEof

服务端运行时对出站的每次 poll 都使用其唤醒注册表中该 key 专属的 waker。客户端运行时把这个 Context 往下传给上游 stream,因此上游就绪时恰好唤醒那个 key。

除往返之外,该测试还固定了三点:

  • 上游管道只能容纳 4 字节,而头部是 6 字节(\xC0\x04dst1),所以在上游读取之前头部无法全部发出。在此期间客户端看不到 +。
  • 客户端 codec 编码的是服务端协议核心的目标(dst1),而不是上游的地址(upstream)。
  • 服务端运行时的 Traffic 统计的是穿过出站明文侧的字节,而不是上游线路上的字节:outbound_tx == 23(仅载荷,不含头部和分帧)和 outbound_rx == 8。
不变量 由谁保证 由谁固定
握手 Done 之前不封装、不交付任何明文 每个明文侧 poll 开头的 ready();poll_read 和 poll_recv_from 中的 is_ready() 守卫 plaintext_waits_for_every_handshake_reply
请求头先于第一帧发出 staging 是先进先出的;start 在任何 seal 之前写入头部 handshake_then_seal_and_open_over_the_wire
拨号 resolve 意味着已拨通、已握手、已 flush ProxyClientConnecting::poll 调用 poll_connected(ready、drain、poll_flush) server_runtime_relays_through_a_client_runtime_in_one_task
上游拒绝变成 ConnectFailed,而不是 Connected 拨号和握手错误让连接 future 以 Err resolve upstream_dial_failure_is_connect_failed_not_connected、refused_upstream_handshake_is_connect_failed_not_connected
违反 codec 约定会失败,绝不空转 handshake 中的无进展检查和 consumed > len 检查;check_frame;seal 字节数检查 handshake_step_without_progress_is_rejected、zero_length_frame_is_rejected、end_marker_beyond_the_slice_is_rejected、datagram_zero_length_frame_is_rejected
每条流的内存有上限 两个固定大小的 BUF_SIZE 数组;make_room 等待上游而不是扩容 backpressure_from_the_wire_reaches_the_writer
大于缓冲区的上游帧会被拒绝 down.is_saturated() 时 fill 失败 没有直接测试
已打开的帧未交付完时 down 不会被 compact opened 已设置时 poll_read 在 fill 之前返回;fill 中的 debug_assert! small_reads_take_one_frame_across_several_calls 覆盖了这条路径
拨号失败后一直保持失败 Wire::Down 是终态,之后每次调用都返回 NotConnected dial_failure_surfaces_on_first_use
被截断的上游帧是错误,而不是干净的 EOF eof() 在线路 EOF 时检查是否还有未解析字节 truncated_upstream_frame_is_unexpected_eof
一次接收就是一个报文 poll_recv_from 在一个数据帧后返回 datagram_codec_sends_and_receives_packets_over_the_stream
BUF_SIZE 留有封装空间 ProxyClientRuntime::new 中的 assert! 构造时强制

每种失败都以 io::Error 的形式到达调用方:

情况 Kind 信息 之后的线路状态
内层拨号失败 拨号错误的 kind 拨号错误 Down
内层 connector 产出数据报 link Unsupported proxy client runtime needs a stream to the upstream, got a datagram socket Down
Down 之后的任何调用 NotConnected 无 Down
start 失败 InvalidData codec 错误 Down
reply、seal、open、finish、seal_to 或 open_from 失败 InvalidData codec 错误 不变
上游写入或读取失败 I/O 错误的 kind I/O 错误 Down
上游写入返回 0 WriteZero 无 Down
连接期间上游 flush 失败 I/O 错误的 kind I/O 错误 Down
poll_flush 或 poll_shutdown 中上游 flush 或 shutdown 失败 I/O 错误的 kind I/O 错误 不变
上游在握手期间关闭 UnexpectedEof upstream closed during the handshake Down
上游在一帧中间关闭 UnexpectedEof upstream closed inside a frame 不变
流结束后进行数据报接收 UnexpectedEof 无 不变
未解析的上游字节占满 down InvalidData upstream frame larger than the client runtime's buffer 不变
报文即使在空的 staging 中也放不下 InvalidInput frame larger than the client runtime's buffer 不变
codec 拒绝了一个放得下的报文 InvalidInput codec refused a packet that fits its reserve 不变
Step 消费的字节多于给出的 Other codec consumed more reply than it was given 不变
Step 没有取得进展 Other codec handshake step made no progress 不变
open 或 open_from 返回了不可能的帧 Other codec opened an impossible frame 不变
seal 取走 0 个或多于所提供的字节 Other codec sealed an impossible byte count 不变

无论 codec 内部用的是什么 kind,codec 错误的 kind 总是 InvalidData,因为 codec_err 用 io::Error::new(io::ErrorKind::InvalidData, e) 包装它。信息是 codec 自己的。HTTP 的 proxy responded with status 407 保留原文(protocols/tests/pipeline/http.rs 中的 new_server_refuses_bad_credentials_with_407),refused_handshake_fails_the_flow 同时检查 InvalidData kind 和 method refused 文本。

服务端运行时把 poll_write、poll_flush、poll_shutdown、poll_read 或 poll_recv_from 返回的任何错误都视为该出站的致命错误:它 drop 该出站,并以 OutboundError 通知协议核心。poll_send_to 的错误则报告为 SendFailed,key 保持存活。因此在实际中,标记为“不变”的那些行不会留下一个半可用的流出站。

取消。 客户端运行时不拥有 task 或 channel,也不持有自身之外的任何东西。drop 它(因为服务端运行时关闭了该 key、连接已结束,或 ProxyClientConnecting future 在 resolve 之前被 drop)会立即 drop 拨号 future 或上游 stream,以及两个缓冲区。没有需要 join 的东西。

常量 位置 值 含义
BUF_SIZE ProxyClientRuntime 的 const 泛型参数 按协议而定,见下 staging 和 down 的大小。codec 打开的最大上游帧必须能放进 down。
HTTP_BUF app/src/outbound/mod.rs、katana src/outbound/mod.rs 16 * 1024 HTTP CONNECT 客户端
SOCKS_BUF 两者 16 * 1024 SOCKS CONNECT 客户端
TROJAN_BUF app/src/outbound/mod.rs 16 * 1024 Trojan 客户端(仅 app)
VLESS_BUF 两者 16 * 1024 VLESS 客户端
VMESS_BUF 两者 32 * 1024 VMess 客户端
SS_BUF 两者 20 * 1024 Shadowsocks(AEAD)客户端
SS2022_BUF 两者 32 * 1024 Shadowsocks 2022 客户端
HttpConnect::REQUEST_MAX protocols/src/http/codec.rs 1024 start 写入的最大 CONNECT 请求;更长的请求以 http: CONNECT request exceeds the codec's reserve 失败

每条流上,客户端运行时在服务端运行时自身缓冲区之外,还要分配 2 * BUF_SIZE 字节以及装箱的拨号 future。因此一条 VMess 流持有 64 KiB 的客户端侧缓冲区。poll_send_to 能接受的最大报文是 BUF_SIZE - STAGING_RESERVE,因为空的 staging 恰好提供 BUF_SIZE 字节的空间;更大的报文以 frame larger than the client runtime's buffer 失败。

客户端运行时每条上游连接只承载一条流。在一条上游载体上共享多条流的客户端(mux.cool 客户端)需要一条为流做准入的路径,而客户端运行时并不提供。

  1. 以 Target = Destination 和 Error = io::Error 实现 ProxyCoreEncodeHandshake,这是 app 和 katana 中 ProxyClient 要求的约束。在 start 中编码流的目的地。只有上游必须在载荷之前应答时才返回 AwaitReply,并让每个 Step 都消费、写入或结束握手。

  2. 声明如实的 STAGING_RESERVE:你曾写入的最大头部、握手轮次、关闭帧或每次 seal 的开销。这样,对遵守它的 codec 来说,Staging::reserve 和 Staging::put 就不会失败。

  3. 实现 ProxyCoreEncode(每次 seal 最多取一帧,每次 open 打开一帧)或 ProxyCoreEncodeDatagram(每次 seal_to 和每个数据帧对应一个报文)。把响应头当作空帧消费,而不是增加一轮应答。

  4. 选择一个既大于 STAGING_RESERVE、又大于上游可能发送的最大帧的 BUF_SIZE,在 app/src/outbound/mod.rs 中添加一个 *_BUF 常量和一个用 Make 闭包构建的出站变体;如果 katana 需要,也在 katana 的 src/outbound/mod.rs 中做同样的事。

  5. 在 protocols/tests/unit/<protocol>/codec.rs 中添加一个不涉及 I/O、直接驱动 codec 的单元测试,并像现有 codec 那样在 codec 文件中用 #[cfg(test)] #[path = …] mod tests; 引入。再针对真实的服务端协议核心添加一个 new_server_vs_new_client_* pipeline 测试。

concepts/tests/client.rs 中的 concepts 测试用玩具 codec(TinyCodec、TwoRoundCodec、Scripted、TinyUdpCodec)在 tokio::io::duplex 管道上驱动运行时,从而可以单独隔离每种行为。运行方式:cargo test -p etemenanki-concepts --test client。

测试 固定的行为
codec_is_driven_without_any_io codec 是一个普通状态机:start 写入头部,seal 最多取一帧,open 原地解密,遇到不完整的帧返回 NeedMore。
handshake_then_seal_and_open_over_the_wire 头部先于第一帧,长写入被拆成多帧,空帧被跳过,End 即 EOF,shutdown 先发送关闭帧再半关闭线路。
small_reads_take_one_frame_across_several_calls 一个已打开的帧通过每次 3 字节的读取分批交付,且不会被重新打开。
truncated_upstream_frame_is_unexpected_eof 线路 EOF 时 down 中还有不完整的帧,结果为 UnexpectedEof。
dial_failure_surfaces_on_first_use 拨号错误(ConnectionRefused)在第一次写入时出现,之后的调用得到 NotConnected。
backpressure_from_the_wire_reaches_the_writer 上游管道只有 16 字节时,200 字节的 write_all 一直 pending,直到上游读取。
plaintext_waits_for_every_handshake_reply 两轮握手:第一个应答之前只发出问候,最后一个应答之前没有载荷,只写的调用方不会死锁,随最后一个应答发送的数据会被打开。
refused_handshake_fails_the_flow reply 中的 codec 错误以 InvalidData 出现,并带有 codec 的信息。
handshake_step_without_progress_is_rejected 无进展检查在 1 秒超时内触发,而不是空转。
zero_length_frame_is_rejected check_frame 拒绝 consumed == 0。
end_marker_beyond_the_slice_is_rejected check_frame 拒绝消费超出切片的 End。
datagram_zero_length_frame_is_rejected poll_recv_from 上的同一守卫。
datagram_codec_sends_and_receives_packets_over_the_stream seal_to 为报文加上端口后分帧,一次线路写入中的两个报文作为两次接收返回,关闭帧变成 UnexpectedEof。
server_side::server_runtime_relays_through_a_client_runtime_in_one_task 在一个 task 内完成完整中继,头部经由 4 字节管道发出之后才有 Connected,Shutdown 变成 finish,以及明文 Traffic 计数。
server_side::upstream_dial_failure_is_connect_failed_not_connected 被拒绝的拨号以 ConnectFailed 到达服务端协议核心(客户端看到 -,绝不会看到 +)。
server_side::refused_upstream_handshake_is_connect_failed_not_connected 上游握手未完成期间不作任何回复,被拒绝的方法变成 ConnectFailed。

真实的 codec 在 protocols/tests/pipeline/*.rs 中与真实的服务端协议核心对跑,借助 protocols/tests/support/pipeline.rs 中的辅助函数 client,它在普通 TCP 上构建 ProxyClientRuntime。运行方式:cargo test -p etemenanki-protocols --test pipeline。

测试 文件 固定的行为
new_server_vs_new_client_tcp http.rs、socks.rs、trojan.rs、vless.rs、vmess.rs stream codec 与真实服务端协议核心之间 70 000 到 100 000 字节的回显往返,然后 shutdown 并干净地 EOF
new_server_vs_new_client_udp trojan.rs、vless.rs、vmess.rs 通过 DatagramLink 的数据报 codec,包括一个 1500 字节的报文(VMess 为 4000 字节)
ss_new_server_vs_new_client_tcp、ss2022_new_server_vs_new_client_tcp、ss2022_multi_user_new_server_vs_new_client shadowsocks.rs Shadowsocks 和 Shadowsocks 2022 的 stream codec
new_server_refuses_bad_credentials_with_407 http.rs HTTP 407 让第一次写入失败,信息中带有该状态
new_server_refuses_an_unreachable_target_after_trying socks.rs SOCKS 拒绝(rejects request)让第一次写入失败,而不是之后的某次读取

每个 codec 在 protocols/tests/unit/<protocol>/codec.rs 下还有不涉及 I/O 直接调用它的单元测试,做法与 codec_is_driven_without_any_io 相同。codec 文件以 #[cfg(test)] 模块引入它们,所以它们随库测试一起运行(cargo test -p etemenanki-protocols --lib)。例如,anonymous_connect_takes_two_rounds 和 credentials_add_a_round_and_a_refusal_is_an_error 覆盖 SOCKS,stream_codec_consumes_the_response_header_as_an_empty_frame 覆盖 VLESS,datagram_codec_refuses_a_packet_that_does_not_fit 覆盖 VMess。