跳转到内容

传输层:WebSocket 与 gRPC

源码文件:22 个 · 核对版本 Etemenanki 596916d · katana v3.0.1
  • Etemenanki/protocols/src/transports/ws/endpoint.rs
  • Etemenanki/protocols/src/transports/ws/stream.rs
  • Etemenanki/protocols/src/transports/grpc/framing.rs
  • Etemenanki/protocols/src/transports/grpc/liveness.rs
  • Etemenanki/protocols/src/transports/grpc/settings.rs
  • Etemenanki/protocols/src/transports/grpc/stream.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/transports/stream.rs
  • Etemenanki/app/src/transport.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/serve.rs
  • Etemenanki/protocols/tests/unit/transports/ws_endpoint.rs
  • Etemenanki/protocols/tests/unit/transports/grpc_framing.rs
  • Etemenanki/protocols/tests/unit/transports/grpc_liveness.rs
  • Etemenanki/protocols/tests/unit/transports/accept.rs
  • Etemenanki/protocols/tests/pipeline/transports.rs
  • Etemenanki/app/tests/integration/e2e_xray.rs
  • Etemenanki/app/tests/integration/e2e_xray_vmess.rs
  • Etemenanki/app/tests/integration/e2e_xray_mux.rs
  • katana/src/inbound.rs

WebSocket 和 gRPC 是让代理连接看起来像 Web 流量的两种承载。两者都位于 etemenanki-protocols 的 protocols/src/transports/ 下,最终落到同一处:一个实现了 AsyncRead + AsyncWrite 的 TransportStream 变体,因此协议核心(core)永远不知道自己运行在哪种承载之上。

本页面向修改承载本身的贡献者:WebSocket 的升级和 early data 处理,gRPC 的 Hunk/MultiHunk 分帧,一个被服务的 socket 所变成的 HTTP/2 连接,以及回收失联对端的定时器。TCP、TLS、MaybeTlsStream 和 TCP keepalive 见传输层:TCP 与 TLS。

组件 文件 → 符号 负责 交给其他部分
WebSocket 端点 ws/endpoint.rs → WsRoute、WsTarget 路径规范化、?ed= 解析、Host 检查、early data 编码、tungstenite 限制 TLS(在其下方)、payload(在其上方)
WebSocket 流 ws/stream.rs → WsStream Binary 消息与字节互转、Ping/Pong、把 Close 视为 EOF、keepalive Ping、空闲拆除、客户端延迟升级 payload 内部的协议分帧
gRPC 分帧 grpc/framing.rs → encode_hunk、encode_multi_hunk、HunkDecoder gRPC 长度前缀与 protobuf data 字段、消息大小上限 HTTP/2 分帧(h2 crate)
gRPC 设置 grpc/settings.rs → GrpcPaths、GrpcMode h2 builder 设置、/<service>/Tun 与 /<service>/TunMulti 路径、请求头与响应头、流的关闭
gRPC 流 grpc/stream.rs → GrpcStream 把一条 HTTP/2 流当作字节流、受流量控制约束的写入、窗口额度、客户端连接驱动 服务端的连接级监管
服务 accept.rs → InboundTransport::accept、serve_h2 10 秒的传输层握手时限、把一个 HTTP/2 连接展开为多条流 逐条流的服务,由调用方的 sink 负责
存活检测 grpc/liveness.rs → Liveness 被服务的 HTTP/2 连接的空闲截止时间与 PING/PONG
拨号 connect.rs → TransportKind、TransportConnector 解析、连接、keepalive,再包上 TLS、WebSocket 或 gRPC

每个被接受或拨出的 socket 都先设置 TCP keepalive,然后是可选的 TLS,最后是承载:

flowchart LR
  tcp["TcpStream"] --> ka["set_keepalive"]
  ka --> tls["MaybeTlsStream"]
  tls --> ws["WsStream"]
  tls --> h2["h2 Connection"]
  h2 --> g1["GrpcStream"]
  h2 --> g2["GrpcStream"]
  ws --> ts["TransportStream"]
  g1 --> ts
  g2 --> ts
  ts --> core["协议核心"]

一个被接受的 WebSocket socket 恰好产生一条流。一个被接受的 gRPC socket 则在连接存续期间,对端每打开一条 HTTP/2 流就产生一条流。

调用方通过入站和出站两个枚举使用这两种承载;WsStream::accept、WsStream::connect 和 GrpcStream::connect 也是公开的,但应用和 katana 只使用这两个枚举及其辅助构造函数。

protocols/src/transports/accept.rs
pub const TRANSPORT_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
pub type Accepted = TransportStream;
pub enum InboundTransport {
Tcp,
Tls(ServerConfig),
Ws {
route: WsRoute,
tls: Option<ServerConfig>,
},
Grpc {
paths: Arc<GrpcPaths>,
tls: Option<ServerConfig>,
},
}
impl InboundTransport {
pub fn ws(path: impl AsRef<str>, host: Option<&str>, tls: Option<ServerConfig>) -> Self;
pub fn grpc(service: impl AsRef<str>, tls: Option<ServerConfig>) -> Self;
pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);
}
protocols/src/transports/connect.rs
pub enum TransportKind {
Tcp,
Tls(ClientConfig),
Ws {
target: WsTarget,
tls: Option<ClientConfig>,
},
Grpc {
authority: Arc<str>,
service: Arc<str>,
mode: GrpcMode,
user_agent: Option<Arc<str>>,
tls: Option<ClientConfig>,
},
}
impl TransportKind {
pub fn ws(host: impl AsRef<str>, path: impl AsRef<str>, tls: Option<ClientConfig>) -> Self;
pub fn grpc(
authority: impl AsRef<str>,
service: impl AsRef<str>,
tls: Option<ClientConfig>,
) -> Self;
pub fn multi(mut self) -> Self;
pub fn user_agent(mut self, agent: Option<&str>) -> Self;
}

对 WebSocket 而言,InboundTransport::accept 在把唯一一条流交给 sink 后立即返回。对 gRPC 而言,它要等 HTTP/2 连接结束才返回,因为这个 future 本身就是连接驱动(见服务端处理 HTTP/2 连接)。因此调用方要在独立的 task 中运行 accept,并为产出的每条流各 spawn 一个 task;应用在 serve_socket 中这样做,详见服务入站。

设置 etemenanki-app 入站 etemenanki-app 出站 katana 入站
WebSocket 路径 ws.path,默认 "/" ws.path,默认 "/" 节点路径,为空时取 "/"
WebSocket Host ws.host,否则不检查 依次取 ws.host、tls.server_name、server;报错 ws stream needs ws.host or server 节点 host,为空时不检查
WebSocket TLS ALPN Alpn::Http1 Alpn::Http1 Alpn::None
gRPC service grpc.service_name,必填 grpc.service_name,必填 节点 service name
gRPC :authority 不检查 依次取 grpc.authority、tls.server_name、server;报错 grpc stream needs grpc.authority or server 不检查
gRPC TLS ALPN Alpn::Http2 Alpn::Http2 Alpn::Http2
gRPC 模式与 user agent 两条路径都服务 始终为 GrpcMode::Gun、DEFAULT_USER_AGENT 两条路径都服务

缺少 grpc.service_name 时,app/src/transport.rs → resolve_stream 校验失败,报错 grpc stream needs grpc.service_name。库中有 TransportKind::multi 和 TransportKind::user_agent,但应用不会调用它们。面向用户的配置键见传输层。

protocols/src/transports/ws/endpoint.rs
pub(crate) const MAX_EARLY_DATA: usize = 16 * 1024;
pub(crate) const MAX_WS_MESSAGE_LEN: usize = 1024 * 1024;
pub(crate) fn ws_config() -> tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
pub(crate) type EarlyDataSlot = Arc<Mutex<Option<Bytes>>>;
pub(crate) fn normalize_path(path: &str) -> String;
pub(crate) fn parse_early_data_path(path: &str) -> (String, Option<usize>);
pub(crate) fn host_matches(request_host: &str, config: &str) -> bool;
pub(crate) fn decode_early_data_header(value: &str) -> Result<Option<Bytes>, ()>;
pub(crate) fn encode_early_data_header(bytes: &[u8]) -> String;
pub struct WsRoute {
path: Arc<str>,
host: Option<Arc<str>>,
}
impl WsRoute {
pub fn new(path: impl AsRef<str>) -> Self;
pub fn host(mut self, host: impl Into<Arc<str>>) -> Self;
pub(crate) fn callback(&self, early_data: EarlyDataSlot) -> WsCallback;
fn accepts(&self, path: &str, host: Option<&str>) -> bool;
}
pub struct WsTarget {
host: Arc<str>,
path: Arc<str>,
early_data_limit: Option<usize>,
}
impl WsTarget {
pub fn new(host: impl AsRef<str>, path: impl AsRef<str>) -> Self;
pub fn early_data_limit(&self) -> Option<usize>;
pub(crate) fn client_request(
&self,
tls: bool,
early_data: Option<&Bytes>,
) -> std::io::Result<http::Request<()>>;
}

WsRoute 描述入站接受什么;WsTarget 描述出站请求什么。WsCallback 实现了 tungstenite 的 Callback,在服务端握手过程中运行,持有路由的一份克隆和共享的 EarlyDataSlot。slot 中的 parking-lot Mutex 每个连接最多加锁两次:一次是回调存入合法的 early data,一次是 WsStream::accept 取出它。

protocols/src/transports/ws/stream.rs
pub const WS_IDLE_TIMEOUT: Duration = Duration::from_secs(300);
pub const WS_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(60);
pub struct WsStream<S> {
state: State<S>,
pending_read: Bytes,
read_closed: bool,
close_sent: bool,
control: Option<Message>,
keepalive: Pin<Box<Sleep>>,
idle: Pin<Box<Sleep>>,
read_waker: Option<Waker>,
}
impl<S> WsStream<S>
where
S: AsyncRead + AsyncWrite + Unpin,
{
pub fn from_upgraded(ws: WebSocketStream<S>, initial: Option<Bytes>) -> Self;
pub async fn accept(stream: S, route: &WsRoute) -> io::Result<Self>;
pub async fn connect(stream: S, target: &WsTarget, secure: bool) -> io::Result<Self>
where
S: Send + 'static;
pub fn is_open(&self) -> bool;
}
impl<S> AsyncRead for WsStream<S> where S: AsyncRead + AsyncWrite + Unpin { /* … */ }
impl<S> AsyncWrite for WsStream<S> where S: AsyncRead + AsyncWrite + Unpin + Send + 'static { /* … */ }

在 TransportStream 中,该流总是装箱后的 WsStream<MaybeTlsStream>。

WsRoute::new 把配置的路径交给 normalize_path:它先用 parse_early_data_path 去掉 ?ed=,再按 Xray 的 GetNormalizedPath 处理(空路径变为 /,缺少开头的 / 时补上)。WsTarget::new 自己调用 parse_early_data_path 以保留上限值,再调用 normalize_path。因此两端可以写同一个路径字符串:服务端和客户端一样会去掉 ?ed=N,线上也永远看不到它。

配置的路径 升级路径 early data 上限(WsTarget)
"" / 无
echo /echo 无
/x?ed=2048 /x Some(2048)
x?ed=2048 /x Some(2048)
/x?foo=1&ed=2048 /x?foo=1 Some(2048)
/x?foo=1 /x?foo=1 无

parse_early_data_path 的规则:

  • 只有名为 ed 且值非空的查询参数才被视为 early data 配置。无论其值能否解析,它都会从路径中移除。
  • 值按 usize 解析。无法解析的值同样会被移除,但上限保持为 None。
  • 存在 ed 时,空的查询参数被丢弃,其余参数用 & 重新拼接。没有 ed 时,路径原样返回。
  • WsStream::connect 只在上限大于零时才延迟升级,所以 ?ed=0 表示不使用 early data。

host_matches 对应 Xray 的 internet.IsValidHTTPHost:两边都转为小写,请求中的值在最后一个 : 处截断以去掉端口,得到的 host 部分必须与配置值完全相等。

路由 请求 结果
未配置 host 任意 Host,或没有 接受
example.com Example.COM:443 接受
example.com other.example:443 404
example.com:443 example.com:443 404:端口只从请求值中去掉,配置值中的端口不会去掉
example.com 没有 Host 头,或值不是可见 ASCII 404
任意 路径与路由不同(区分大小写比较) 404

拒绝响应由 reject_404 构造:回调返回一个空 body 的 404 Not Found,使 tungstenite 作出应答并让握手失败。

与 Xray 兼容的 early data 把客户端的最初几个字节放进升级请求的 Sec-WebSocket-Protocol 头中,让客户端先发的协议省下一个往返。

步骤 端 机制
编码 客户端 encode_early_data_header:URL 安全的 base64 字母表,无填充(b"hello" → aGVsbG8,[0xfb, 0xff] → -_8)
大小 客户端 取第一次非空写入的前 min(len, ed, MAX_EARLY_DATA) 字节,因此不会超过 16 KiB
规范化 服务端 去掉首尾空白,把 + 映射为 -、/ 映射为 _,删除 =:标准和 URL 安全字母表、有无填充都能解码
解码 服务端 URL_SAFE_NO_PAD。空值、无法解码的值、解码后为空的值都不算 early data
上限 服务端 解码后超过 MAX_EARLY_DATA 字节时返回 Err(()),回调应答 413 Payload Too Large
回显 服务端 合法的 early data 存入 slot,请求中的 Sec-WebSocket-Protocol 值原样回显在 101 响应中
交付 服务端 WsStream::accept 从 slot 取出数据并预置到 pending_read,因此 early data 是协议核心读到的第一段数据

不是 early data 的头会让 slot 保持为空,响应中也不带 Sec-WebSocket-Protocol;升级照常成功。

回显是必需的。如果请求带了 Sec-WebSocket-Protocol,而 101 响应中没有,tungstenite 的客户端握手会将其视为子协议错误并让升级失败;所以一个消费了 early data 却不回显的服务端,会让每个使用 ?ed= 的 Etemenanki 客户端都无法连接。

在客户端,上限为正数时 WsStream::connect 根本不做升级。它把 socket 暂存在 State::Deferred 中并返回。第一次非空的 poll_write 调用 start_upgrade:取出至多 limit 字节,用它们构造请求,把 tungstenite 握手装箱放进 State::Upgrading,唤醒停在 read_waker 上的读者,并返回 Ok(take)。装箱的握手在流再次被 poll 之前不做任何事:write_all 的剩余部分、被唤醒的读、flush 或 shutdown 都会通过 poll_open 推动它。升级完成后,调用方的 write_all 再把剩余数据作为普通消息写出。

sequenceDiagram
  participant CC as 客户端核心
  participant CW as WsStream 客户端
  participant SW as WsStream::accept
  participant SC as 服务端核心
  CC->>CW: connect,路径 /ws?ed=2048
  Note over CW: State Deferred,未发送任何字节
  CC->>CW: write_all 2100 字节,第一次 poll_write
  CW->>CW: start_upgrade 取 2048 字节,状态 Upgrading
  CW-->>CC: Ok(2048)
  CC->>CW: 第二次 poll_write,52 字节
  Note over CW: poll_open 推动装箱的握手
  CW->>SW: GET /ws,Sec-WebSocket-Protocol 为 base64url(2048 字节)
  SW->>SW: WsRoute::accepts 检查路径和 Host
  SW->>SW: decode_early_data_header 存入 EarlyDataSlot
  SW-->>CW: 101,回显 Sec-WebSocket-Protocol
  Note over CW: Upgrading 变为 Open
  SW->>SC: WsStream,pending_read = 2048 字节
  CW->>SW: Binary 消息,52 字节
  SC->>SC: 先读到 2048 字节,再读到 52 字节

在第一次写入之前发起的读会从 poll_open 得到 Pending 并保存其 waker,由 start_upgrade 唤醒。处于 Deferred 时调用 poll_shutdown 会以无 early data 的方式开始升级,因此关闭仍会以 WebSocket Close 帧的形式送达对端。处于 Deferred 时的 poll_flush 什么也不做,直接返回 Ok(())。

stateDiagram-v2
  [*] --> Open: accept,或不带 ed 的 connect
  [*] --> Deferred: ed 大于 0 的 connect
  Deferred --> Upgrading: 第一次非空 poll_write,或 poll_shutdown
  Upgrading --> Open: 握手完成,重置定时器
  Upgrading --> Failed: 握手出错
  Failed --> Failed: 每次调用都返回保存的错误
  Open --> [*]

State::Failed(io::ErrorKind, String) 保存升级错误的种类和文本,之后每次调用 poll_open 都会重建一个等价的 io::Error。触发升级的那次写入已经返回了 Ok,所以失败会在下一次读、写、flush 或 shutdown 时暴露。在 Upgrading 状态下丢弃流,会连同装箱的握手 future 和 socket 一起丢弃。

start_upgrade 在构造请求之前就把 socket 从 Deferred 中取出。如果此时 WsTarget::client_request 失败(例如 host 无法构成合法的 URI),第一次写入会返回该错误,状态停留在没有 socket 的 Deferred,之后的任何写入或 shutdown 都以 websocket upgrade already started 失败。flush 仍返回 Ok(()),而读会一直停在 Pending,因为已经没有任何东西能唤醒它。

WebSocket 消息与字节之间的映射是固定的:

方向 事件 WsStream 的动作
读 Binary(bytes) 存入 pending_read;按缓冲区大小分多次 poll_read 交出
读 Ping(payload) control = Some(Pong(payload)),由下一次 poll_control 发送
读 Close(_) read_closed = true:EOF
读 Text、Pong、原始 Frame 忽略,但仍会重置两个定时器
读 流结束、ConnectionClosed、AlreadyClosed、ResetWithoutClosingHandshake EOF
读 其他任何 tungstenite 错误,包括大小超限 Err;若不是 I/O 错误则为 io::Error::other
写 空缓冲区 Ok(0),不发消息
写 非空缓冲区 恰好一条携带缓冲区副本的 Binary 消息,然后返回 Ok(buf.len())
flush 发送等待中的控制帧,然后 flush sink
shutdown 发送一次 Close(None)(由 close_sent 保证只发一次),然后 flush;flush 期间遇到 ConnectionClosed 或 AlreadyClosed 视为成功

如果写入时 sink 报告 ConnectionClosed 或 AlreadyClosed,写入以 BrokenPipe 失败,文本为 websocket closed(ws_err)。

control 最多保存一个控制帧。poll_control 在 poll_read 和写路径中都会被调用;在读路径上,sink 未就绪不会阻塞读:控制帧继续等待,读照常进行。

tungstenite 0.30 读到 Ping 时也会自行排队一个 Pong,而交给其 sink 的 Pong 会替换尚未 flush 的已排队 Pong。因此对端每个 Ping 看到一个 Pong;如果 tungstenite 在 WsStream 发出第二个 Pong 之前已经 flush 了自己的 Pong,则会看到两个。两个 Pong 都携带 Ping 的 payload,而 RFC 6455 允许主动发送的 Pong。

tungstenite 自身不发送 keepalive,所以 WsStream 自己持有两个 tokio::time::Sleep 定时器:

定时器 常量 触发时机 效果
keepalive WS_KEEPALIVE_INTERVAL = 60 秒 最后一次活动后 60 秒 重新设定 60 秒;若没有等待中的控制帧,则排队一个空 payload 的 Ping
idle WS_IDLE_TIMEOUT = 300 秒 最后一次活动后 300 秒 read_closed = true:读端报告 EOF,就像对端没发 FIN 就消失了一样

touch 同时重置两者。它在升级完成时、每收到任何类型的帧时(包括 Pong,正是这一点让 Ping 成为存活探测),以及每次写入成功后运行。两个定时器都在 poll_read 中被 poll,所以只要有读者在等待该流,它们就会触发。这些值与 HTTP/2 的空闲和 keepalive 常量一致,使两种承载的行为相同。

ws_config 用 MAX_WS_MESSAGE_LEN = 1 MiB 替换 tungstenite 的默认值(每条消息 64 MiB,每帧 16 MiB),同时用于 max_message_size 和 max_frame_size。tungstenite 会先把分片的消息重组到一个缓冲区再交出,所以这就是单个对端能让一个会话为一条消息占用的最大内存。同一份配置传给 accept_hdr_async_with_config 和 client_async_with_config,因此对入站会话和拨出的会话同样适用,并与 gRPC 承载的 MAX_GRPC_MESSAGE_LEN 一致。

tungstenite 的这两个限制都只作用于接收的消息。WsStream::poll_write 从不拆分缓冲区,所以一次超过 1 MiB 的写入会作为一条消息发出,而 Etemenanki 对端会以大小错误拒绝它。协议核心写入的块远小于此,pipeline 测试一次写入 200 000 字节。

gRPC 承载就是 Xray 的 “gun” 传输:一个双向流式 gRPC 调用,其消息是 protobuf Hunk { bytes data = 1; },在 multi 模式下则是 MultiHunk { repeated bytes data = 1; }。这里没有生成的 protobuf 代码,唯一的字段由手工写入和解析。

protocols/src/transports/grpc/settings.rs
pub struct GrpcPaths {
tun: Arc<str>,
multi: Arc<str>,
}
impl GrpcPaths {
pub fn new(service: impl AsRef<str>) -> Self;
pub fn classify(&self, path: &str) -> Option<GrpcMode>;
}
pub enum GrpcMode {
Gun,
Multi,
}
pub(crate) fn configured_client_builder() -> h2::client::Builder;
pub(crate) fn configured_server_builder() -> h2::server::Builder;
pub(crate) fn grpc_response() -> http::Response<()>;
pub(crate) fn grpc_request(
secure: bool,
authority: &str,
service: &str,
mode: GrpcMode,
user_agent: Option<&str>,
) -> std::io::Result<http::Request<()>>;
pub(crate) fn finish_stream(send_stream: &mut SendStream<Bytes>, is_server: bool);
protocols/src/transports/grpc/framing.rs
pub(crate) fn encode_hunk(payload: &[u8]) -> io::Result<Bytes>;
pub(crate) fn encode_multi_hunk(payloads: &[Bytes]) -> io::Result<Bytes>;
pub(crate) struct HunkDecoder {
chunks: VecDeque<Bytes>,
len: usize,
}
impl HunkDecoder {
pub(crate) fn new() -> Self;
pub(crate) fn push(&mut self, chunk: Bytes);
pub(crate) fn buffered(&self) -> usize;
pub(crate) fn next(&mut self) -> io::Result<Option<Bytes>>;
pub(crate) fn next_multi(&mut self) -> io::Result<Option<Vec<Bytes>>>;
}
protocols/src/transports/grpc/stream.rs
pub struct GrpcStream {
send: SendStream<Bytes>,
recv: Recv,
mode: GrpcMode,
is_server: bool,
decoder: HunkDecoder,
ready: VecDeque<Bytes>,
write_pending: Option<Bytes>,
finished: bool,
_driver: Option<AbortOnDropHandle<()>>,
_guard: Option<StreamGuard>,
}
enum Recv {
Awaiting(ResponseFuture),
Body(RecvStream),
Done,
}
impl GrpcStream {
pub(crate) fn served(
send: SendStream<Bytes>,
recv: RecvStream,
mode: GrpcMode,
count: &Arc<StreamCount>,
) -> Self;
pub async fn connect<T>(
io: T,
secure: bool,
authority: &str,
service: &str,
mode: GrpcMode,
user_agent: Option<&str>,
) -> io::Result<Self>
where
T: AsyncRead + AsyncWrite + Unpin + Send + 'static;
}
pub(crate) struct StreamCount {
live: AtomicUsize,
changed: Notify,
}

被服务的 GrpcStream 带有 StreamGuard、没有 driver;拨出的 GrpcStream 带有 driver、没有 guard。

GrpcPaths::new 用 service name 构造两条路径,名字原样插入,不做任何规范化。classify 对请求路径做精确比较。

路径 GrpcMode 消息类型
/<service>/Tun Gun 每条 gRPC 消息一个 Hunk
/<service>/TunMulti Multi 每条 gRPC 消息一个 MultiHunk
其他任何路径 无 以 REFUSED_STREAM 重置该流
消息 头部
客户端请求(grpc_request) POST,URI 为 https://<authority>/<service>/Tun(或 http://,或 /TunMulti),content-type: application/grpc,te: trailers,设置了 user agent 时带 user-agent
服务端响应(grpc_response) 200,content-type: application/grpc,流一被接受就发送
服务端关闭(finish_stream) trailers grpc-status: 0
客户端关闭(finish_stream) 一个带 END_STREAM 的空 DATA 帧

DEFAULT_USER_AGENT 是一个桌面版 Chrome 字符串:Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/152.0.7977.42 Safari/537.36。TransportKind::grpc 会设置它;TransportKind::user_agent(None) 会去掉该头。

流上的每条 gRPC 消息都是带长度前缀的消息:

偏移 大小 字段 含义
0 1 压缩标志 写为 0x00:未压缩。解码器不读取这个字节
1 4 消息长度 u32,大端序:其后 protobuf 消息的长度
5 消息长度 消息 一个 protobuf Hunk 或 MultiHunk

Hunk 消息是一个长度分隔字段:

大小 字段 含义
1 tag 0x0A(HUNK_DATA_TAG):字段号 1,wire type 2
1 到 10 长度 base-128 varint:data 的长度
长度 data 隧道中传输的字节

MultiHunk 消息是零个或多个同样的 0x0A / varint / data 三元组首尾相接。encode_multi_hunk 跳过空 payload,消息长度只按非空的部分计算。

测试中的两个编码,逐字节列出:

调用 字节
encode_hunk(b"ok") 00 00 00 00 04 0A 02 6F 6B
encode_multi_hunk(["a", "", "bc"]) 00 00 00 00 07 0A 01 61 0A 02 62 63

130 字节的 payload 需要两字节的 varint(82 01),因此其消息长度为 0x85(133)。

GrpcStream::poll_write 把每次非空写入编码为一条消息:Gun 模式下是一个 Hunk,Multi 模式下是只有一个条目的 MultiHunk。只有当消息长度放不进 u32 时,编码才会以 hunk message too large 失败。编码器不应用 MAX_GRPC_MESSAGE_LEN,所以与 WebSocket 一样,一次超过约 1 MiB 的写入会产生一条被 Etemenanki 对端解码器拒绝的消息。

HTTP/2 DATA 帧与 gRPC 消息并不对齐,所以 HunkDecoder 把收到的块保存在一个 VecDeque<Bytes> 中并维护累计长度,在凑齐一整条消息之前从不复制:

  1. take_frame 跨块边界窥视 5 字节的头部(copy_prefix::<5>)。不足 5 字节:Ok(None)。
  2. 从第 1 到 4 字节读出长度,超过 MAX_GRPC_MESSAGE_LEN 就立即拒绝(此时 body 尚未到达),报错 grpc message length exceeds the accepted maximum。
  3. 已缓冲的字节少于 5 + length 时返回 Ok(None) 并等待。
  4. 否则丢弃头部并切出 body。body 位于单个块中时,这是零复制的 split_to;否则把各段复制到一个 BytesMut 中。
  5. next(gun)要求第一个字节是 0x0A,读取 varint 并切出 data。空消息或空 data 会被跳过,循环继续处理下一条消息。同一消息中声明的 data 之后的字节会被丢弃。
  6. next_multi 遍历消息的每个字段,要求每个都是 0x0A,并返回非空的条目。它可能返回空的 Vec。
条件 错误文本(io::ErrorKind::InvalidData)
声明长度超过 1 MiB grpc message length exceeds the accepted maximum
长度加头部溢出 usize grpc frame too large
字段 tag 不是 0x0A unexpected protobuf field in Hunk / unexpected protobuf field in MultiHunk
varint 越过末尾 truncated varint
varint 在第十个字节之后仍未结束 varint overflow
varint 放不进 usize hunk data too large
字节数少于 varint 声明的长度 hunk data shorter than declared
常量 值 客户端 builder 服务端 builder
H2_INITIAL_STREAM_WINDOW_SIZE 4 MiB 是 是
H2_INITIAL_CONNECTION_WINDOW_SIZE 16 MiB 是 是
H2_MAX_FRAME_SIZE 256 KiB 是 是
H2_MAX_CONCURRENT_STREAMS 256 否 是
MAX_GRPC_MESSAGE_LEN 1 MiB 解码器 解码器
H2_IDLE_TIMEOUT 300 秒 否 Liveness
H2_KEEPALIVE_INTERVAL 60 秒 否 Liveness
H2_KEEPALIVE_TIMEOUT 20 秒 否 Liveness
DEFAULT_USER_AGENT Chrome 字符串 TransportKind::grpc 否

窗口大小是为高 RTT 路径上的隧道吞吐量而设定的。h2 builder 的其余设置都是 h2 crate 的默认值。

sequenceDiagram
  participant CC as 客户端核心
  participant CG as GrpcStream 客户端
  participant D as driver task
  participant S as serve_h2
  participant SG as 被服务的 GrpcStream
  participant SC as 服务端核心
  CG->>D: 握手,然后 tokio::spawn(conn)
  CG->>S: HEADERS POST /svc/Tun
  S->>S: GrpcPaths::classify 得到 Gun
  S-->>CG: HEADERS 200 application/grpc
  S->>SG: GrpcStream::served,取得 guard
  S->>SC: sink(TransportStream::Grpc)
  CC->>CG: poll_write payload
  CG->>SG: DATA encode_hunk(payload)
  SG->>SC: HunkDecoder::next,release_capacity
  SC->>SG: poll_write 应答
  SG->>CG: DATA encode_hunk(reply)
  CG->>CC: Recv 从 Awaiting 变为 Body,解码
  CC->>CG: poll_shutdown
  CG->>SG: 带 END_STREAM 的空 DATA
  SC->>SG: poll_shutdown
  SG->>CG: trailers grpc-status 0

客户端在 connect 中不等待响应头。Recv::Awaiting 持有 ResponseFuture,第一次 poll_read 把它解析为 Recv::Body,所以客户端先发的协议可以在 connect 返回后立即写入。

对于 Grpc,InboundTransport::accept 在 10 秒的 within("grpc", …) 时限内完成(可选的)TLS 和 h2 服务端握手,然后调用 serve_h2:

protocols/src/transports/accept.rs
async fn serve_h2<T, F>(
mut conn: Connection<T, Bytes>,
paths: &GrpcPaths,
sink: &mut F,
) -> io::Result<()>
where
T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
F: FnMut(Accepted),

h2::server::Connection::accept 既产出新流,也驱动连接的 I/O。被服务的流运行在其他 task 中,但只有 serve_h2 持续 poll conn.accept(),它们的字节才会流动。这个循环是一个包含三个分支的 tokio::select!:

flowchart TB
  top["loop:idle = count.is_idle()"] --> sel{"select!"}
  sel -->|"conn.accept() 产出一条流"| cls{"classify 路径"}
  cls -->|"None"| rst["send_reset REFUSED_STREAM"]
  cls -->|"Gun 或 Multi"| resp["send_response 200"]
  resp --> sink["sink(GrpcStream::served)"]
  sel -->|"count.changed(),仅在非 idle 时"| prog["liveness.note_progress()"]
  sel -->|"liveness.watch(idle) 完成"| stop["break:丢弃连接"]
  sel -->|"conn.accept() 产出 None"| stop
  rst --> top
  sink --> top
  prog --> top
  • 每条被接受的流,无论是否被拒绝,都会调用 note_progress,重新开始空闲截止时间。
  • 路径被 classify 拒绝的流以 REFUSED_STREAM 重置;连接继续服务其他流。
  • 如果 send_response 失败,该流被丢弃,循环继续。
  • conn.accept() 出错时,serve_h2 以该错误结束(io::Error::other)。
  • StreamCount 跟踪存活的被服务流。每个 GrpcStream::served 取得一个使 live 加一的 StreamGuard;它的 Drop 使 live 减一并调用 notify_waiters,在 count.changed() 分支等待时唤醒它。

serve_h2 返回时 Connection 被丢弃,该连接上仍打开的每条流都会在下一次读或写时失败。

Liveness 为被服务的连接防范两种风险,各有自己的截止时间:

protocols/src/transports/grpc/liveness.rs
pub(crate) struct Liveness {
ping_pong: Option<PingPong>,
next_ping: Instant,
pong_due: Option<Instant>,
idle_due: Instant,
}
impl Liveness {
pub(crate) fn new(ping_pong: Option<PingPong>) -> Self;
pub(crate) fn note_progress(&mut self);
pub(crate) async fn watch(&mut self, idle: bool);
}
pub(crate) enum Verdict {
Alive,
Dead,
}
风险 截止时间 适用条件 重置方式
对端完成握手却不打开任何流,会让 Connection::accept 永远挂起 idle_due,H2_IDLE_TIMEOUT = 300 秒 idle 为 true:StreamCount 为零 note_progress:有流被接受或有流结束
对端在流仍打开时不发 FIN 就消失,此时空闲截止时间永远不适用 pong_due,每次 PING 后 H2_KEEPALIVE_TIMEOUT = 20 秒 有未应答的 PING 对应的 PONG

next_ping 从 serve_h2 启动起每隔 H2_KEEPALIVE_INTERVAL(60 秒)触发一次,与流量无关。触发时若没有未应答的 PING,tick 发送 Ping::opaque(),并在同一步中设置 pong_due。send_ping 失败立即判为 Dead。poll_pong 返回 Ok 时清除 pong_due;返回 Err 则为 Dead。serve_h2 传入 conn.ping_pong(),h2 只会交出它一次;若为 None,poll_pong 永远挂起,只有定时器生效。

watch 循环调用 tick,直到判定为 Dead 才返回,而不是每次 tick 都返回。在 tokio::select! 中,模式不匹配的分支会在本次调用余下的时间里被禁用,所以一个每次 tick 都返回的分支,在第一次发现对端存活后就不会再重新设定定时器。这个循环是可取消安全的,这一点是必需的,因为 serve_h2 的 select! 中其他任何分支完成时都会丢弃它:它的 await 只有一个定时器和一次 poll,跨 await 不改变任何状态,并且 pong_due 与发送 PING 在同一个同步步骤中设置,所以被取消的 watch 永远不会留下无人监视的 PING。

Dead 判定使 serve_h2 以 Ok(()) 结束:放弃一个对端不算错误。

GrpcStream::connect 每次拨号都会新建一个 HTTP/2 连接,并在其上打开一条流:

  1. configured_client_builder().handshake(io) 返回 (SendRequest, Connection)。
  2. Connection future 用 tokio::spawn 启动,并包在 AbortOnDropHandle 中。它以错误结束时,该 task 以 debug 级别记录 grpc connection ended: …。
  3. send_req.ready() 等待流容量,send_request(request, false) 打开流,body 保持打开。
  4. connect 返回时 SendRequest 句柄被丢弃;打开着的流维持连接存活。

driver 句柄保存在 _driver 中。丢弃 GrpcStream 会丢弃该句柄,从而中止 driver task 及其连接。拨出的连接从不在流之间共享,也不运行 Liveness。

write_pending: Option<Bytes> 最多保存一条已编码的消息:

  1. poll_write 先用 poll_drain 排空上一条消息。在它完全交给 h2 之前返回 Pending,HTTP/2 流量控制正是以此向写入方施加背压。
  2. 编码新的缓冲区,存入 write_pending,并尽力执行一次 poll_drain。
  3. 即使消息还有一部分未发出,也返回 Ok(buf.len())。剩余部分由 poll_flush 和下一次 poll_write 完成。

poll_drain 调用 reserve_capacity(pending.len()),等待 poll_capacity,再用 send_data(chunk, false) 按获得的容量发送尽可能多的字节。poll_capacity 返回 None 表示流已关闭:BrokenPipe,h2 send stream closed。poll_shutdown 先排空,再调用一次 finish_stream(由 finished 保证只调用一次)。

poll_read 按固定顺序工作:

  1. 如果 ready 中有数据,从其头部交出字节。
  2. 否则运行 decode,向解码器要一条消息(或一批 MultiHunk),把其 payload 追加到 ready。
  3. 否则从 h2 拉取一个 DATA 块(poll_data)送入解码器,或先解析 Recv::Awaiting。poll_data 返回 None 时转入 Recv::Done,即 EOF。

解码器每次只接收一个 DATA 块,且仅在 ready 为空、缓冲中没有完整消息时才接收,所以在读者请求之前不会从 h2 拉取任何数据。

decode 用 flow_control().release_capacity(consumed) 归还窗口额度,其中 consumed 是本次调用前后 HunkDecoder::buffered() 的减少量。未完成消息的字节继续占用其窗口额度,所以在读者消费之前,对端发送的单条消息不能超过流窗口所允许的量;这是在声明长度的 MAX_GRPC_MESSAGE_LEN 上限之外的又一层约束。

不变量 保证机制 固定它的测试
一次非空写入恰好是一条 WebSocket Binary 消息 WsStream::poll_write 发送 Message::Binary(Bytes::copy_from_slice(buf)) 没有专门的测试;pipeline 往返测试只比较回显的字节
一次非空写入恰好是一条 gRPC 消息 GrpcStream::poll_write → encode_hunk / encode_multi_hunk encode_hunk_writes_exact_small_frame、encode_multi_hunk_writes_exact_repeated_fields(protocols/tests/unit/transports/grpc_framing.rs)固定编码;one_grpc_connection_carries_many_streams(protocols/tests/pipeline/transports.rs)检查被服务的流把一次 10 字节的读取恰好回显为一个 Hunk
两端可用同一个路径字符串;?ed= 永远不会上线 WsRoute::new 和 WsTarget::new 中的 parse_early_data_path + normalize_path normalizes_paths、parses_outbound_early_data_path(protocols/tests/unit/transports/ws_endpoint.rs)
随升级携带的最多 min(ed, MAX_EARLY_DATA) 字节,且是服务端核心读到的最初字节 start_upgrade 返回 take;from_upgraded 预置 pending_read ws_early_data_is_the_first_bytes_the_server_reads(protocols/tests/pipeline/transports.rs)
超过 16 KiB 的 early data 在升级前被拒绝 decode_early_data_header 返回 Err(()),回调应答 413 rejects_oversized_early_data_header、callback_rejects_oversized_early_data(ws_endpoint.rs)
合法的 early data 被回显;非法的被忽略且不回显 WsCallback::on_request callback_captures_and_echoes_valid_early_data、callback_ignores_invalid_early_data、decodes_xray_early_data_header、encodes_xray_early_data_header(ws_endpoint.rs)
路径或 Host 不对时永远不会升级 WsRoute::accepts → reject_404 route_accepts_only_matching_path_and_host、host_matching_ignores_case_and_port(ws_endpoint.rs),ws_rejects_a_wrong_path_at_the_upgrade(transports.rs)
不接受任何超过 1 MiB 的 WebSocket 消息或帧,入站和拨出的会话都一样 ws_config 传给 tungstenite 的两个入口 the_configured_limits_replace_tungstenite_defaults(ws_endpoint.rs)检查配置值
声明长度超过 1 MiB 的 gRPC 消息在缓冲其 body 之前就被拒绝 HunkDecoder::take_frame 在头部检查 MAX_GRPC_MESSAGE_LEN decoder_rejects_a_message_longer_than_the_maximum、decoder_rejects_an_oversized_length_before_buffering_the_body、decoder_accepts_a_message_at_the_maximum(grpc_framing.rs)
DATA 帧边界无关紧要 HunkDecoder 的块队列和 copy_prefix decoder_waits_for_split_frame、decoder_reads_multiple_frames_from_one_chunk、decoder_reads_multi_hunk_entries(grpc_framing.rs)
格式错误的 hunk 以 InvalidData 失败 next / next_multi 中的 tag 和长度检查 decoder_rejects_unexpected_field、decoder_rejects_short_declared_payload(grpc_framing.rs)
只服务 /<service>/Tun 和 /<service>/TunMulti GrpcPaths::classify,其他路径 REFUSED_STREAM classify_path_matches_tun_modes(protocols/tests/unit/transports/grpc_liveness.rs)固定分类;没有测试驱动重置
一个被服务的连接承载多条流 serve_h2 循环,每条流调用一次 sink one_grpc_connection_carries_many_streams(transports.rs)
没有流的被服务连接在 H2_IDLE_TIMEOUT 后被丢弃 Liveness::idle_due liveness_gives_up_on_a_connection_with_no_streams、liveness_restarts_the_idle_deadline_on_progress(grpc_liveness.rs),a_connection_that_opens_no_stream_is_given_up_on(protocols/tests/unit/transports/accept.rs)
有存活流的被服务连接不会因空闲截止时间被丢弃 只有 idle 为 true 时才检查 idle_due liveness_keeps_a_connection_carrying_streams(grpc_liveness.rs)
只为已消费的字节归还窗口额度 GrpcStream::decode 释放 held - buffered() 没有专门的测试
拨出的连接随其流一起消亡 _driver 中的 AbortOnDropHandle 没有专门的测试
Liveness::watch 是可取消安全的 pong_due 与 send_ping 在同一步设置 没有专门的测试
情形 结果
入站的 TLS 加 WebSocket 升级,或 TLS 加 h2 握手,超过 TRANSPORT_HANDSHAKE_TIMEOUT accept 以 TimedOut 失败,文本为 websocket handshake timed out 或 grpc handshake timed out;应用以 debug 级别记录 inbound transport failed
升级被拒绝(404、413) 服务端的 accept 失败。不带 ed 拨号的客户端在 connect 中失败,错误文本包含状态码;延迟升级的客户端把错误存入 State::Failed,在下一次调用时报告
延迟升级失败 State::Failed;第一次写入已经返回 Ok(take),之后每次调用都返回保存的错误
WebSocket 对端关闭或消失 读到 EOF(Close、流结束、未经关闭握手的重置)
WebSocket 静默 300 秒 读到 EOF
WebSocket 关闭后写入 BrokenPipe,websocket closed
未知的 gRPC 路径 该流以 REFUSED_STREAM 重置;连接继续
send_response 失败 该流被静默丢弃;连接继续
conn.accept() 出错 serve_h2 返回该错误,连接被丢弃
Liveness 判定对端已死 serve_h2 返回 Ok(()),连接被丢弃,打开着的流失败
GrpcStream 上的 h2 流或连接错误 h2_err:I/O 错误被解包,其他错误变为 io::Error::other
排空时发送端已关闭 BrokenPipe,h2 send stream closed
格式错误的 gRPC 消息 解码器返回 InvalidData,由 poll_read 抛出
拨出的 GrpcStream 被丢弃 driver task 被中止,HTTP/2 连接关闭
拨出的 WsStream 在升级中被丢弃 装箱的握手 future 及其 socket 被丢弃
常量 值 位置
TRANSPORT_HANDSHAKE_TIMEOUT 10 秒 accept.rs:accept 时的 TLS 加升级或 h2 握手
MAX_EARLY_DATA 16 KiB ws/endpoint.rs:解码后的 early data,以及客户端 ed 的上限
MAX_WS_MESSAGE_LEN 1 MiB ws/endpoint.rs:tungstenite 的 max_message_size 和 max_frame_size
WS_KEEPALIVE_INTERVAL 60 秒 ws/stream.rs:发送 Ping 之前的静默时长
WS_IDLE_TIMEOUT 300 秒 ws/stream.rs:报告 EOF 之前的静默时长
MAX_GRPC_MESSAGE_LEN 1 MiB grpc/settings.rs:gRPC 消息的最大声明长度
H2_INITIAL_STREAM_WINDOW_SIZE 4 MiB grpc/settings.rs
H2_INITIAL_CONNECTION_WINDOW_SIZE 16 MiB grpc/settings.rs
H2_MAX_FRAME_SIZE 256 KiB grpc/settings.rs
H2_MAX_CONCURRENT_STREAMS 256 grpc/settings.rs:仅服务端
H2_IDLE_TIMEOUT 300 秒 grpc/settings.rs:没有流的被服务连接
H2_KEEPALIVE_INTERVAL 60 秒 grpc/settings.rs:被服务连接上的 PING 周期
H2_KEEPALIVE_TIMEOUT 20 秒 grpc/settings.rs:PONG 截止时间

TCP keepalive(TCP_KEEPALIVE_IDLE 120 秒,TCP_KEEPALIVE_INTERVAL 30 秒,TCP_KEEPALIVE_RETRIES 3)位于两种承载之下;见传输层:TCP 与 TLS。

文件 测试 覆盖
protocols/tests/unit/transports/ws_endpoint.rs normalizes_paths、parses_outbound_early_data_path、host_matching_ignores_case_and_port、route_accepts_only_matching_path_and_host、decodes_xray_early_data_header、encodes_xray_early_data_header、rejects_oversized_early_data_header、callback_captures_and_echoes_valid_early_data、callback_ignores_invalid_early_data、callback_rejects_oversized_early_data、the_configured_limits_replace_tungstenite_defaults 路径与 ?ed= 处理、Host 检查、early data 编解码与回调、大小限制
protocols/tests/unit/transports/grpc_framing.rs encode_hunk_writes_exact_small_frame、encode_hunk_writes_multi_byte_varint_length、encode_multi_hunk_writes_exact_repeated_fields、decoder_waits_for_split_frame、decoder_reads_multiple_frames_from_one_chunk、decoder_reads_multi_hunk_entries、decoder_rejects_unexpected_field、decoder_rejects_short_declared_payload、decoder_rejects_a_message_longer_than_the_maximum、decoder_rejects_an_oversized_length_before_buffering_the_body、decoder_accepts_a_message_at_the_maximum 精确的线上字节、重组、格式错误的输入、1 MiB 上限
protocols/tests/unit/transports/grpc_liveness.rs classify_path_matches_tun_modes、liveness_gives_up_on_a_connection_with_no_streams、liveness_keeps_a_connection_carrying_streams、liveness_restarts_the_idle_deadline_on_progress 路径分类;在暂停的时钟下测试空闲截止时间
protocols/tests/unit/transports/accept.rs a_connection_that_opens_no_stream_is_given_up_on 通过 tokio::io::duplex 把 serve_h2 和 Liveness 连起来,手动推进时钟
protocols/tests/pipeline/transports.rs ws_round_trip、ws_over_tls_with_early_data_round_trip、ws_early_data_is_the_first_bytes_the_server_reads、ws_rejects_a_wrong_path_at_the_upgrade、grpc_round_trip、grpc_multi_mode_over_tls_round_trip、one_grpc_connection_carries_many_streams 四个 _round_trip 测试通过 InboundTransport 和 TransportConnector 回显 hello 和 200 000 字节,然后期望干净的 EOF。early data 测试用 ?ed=16 写入 19 字节,检查第一次写入返回 16。错误路径测试期望拨号错误中包含 404。多流测试在一个原始 h2 客户端连接上打开三条流
app/tests/integration/e2e_xray.rs app_client_ws_xray_server_plain、app_server_ws_xray_client_plain、app_client_grpc_xray_server_plain、app_server_grpc_xray_client_plain、app_client_ws_xray_server_tls、app_client_grpc_xray_server_tls、app_server_ws_xray_client_tls、app_server_grpc_xray_client_tls 与真实的 Xray 二进制双向互通,VLESS over WebSocket 和 gRPC
app/tests/integration/e2e_xray_vmess.rs app_server_vmess_grpc_xray_client_tls、app_client_vmess_ws_xray_server_early_data_plain、app_server_vmess_ws_xray_client_early_data_plain 与 Xray 互通的 gRPC over TLS 和 WebSocket early data(/vmess?ed=2048)
app/tests/integration/e2e_xray_mux.rs vless_mux_over_ws_tls 在 WebSocket over TLS 之上组合 mux.cool

Xray 互通测试用 go 构建 Xray;没有 go 或构建失败时会自行跳过并打印 SKIP。没有测试驱动 WS_KEEPALIVE_INTERVAL、WS_IDLE_TIMEOUT 或 Liveness 的 PING/PONG 部分;修改这些地方需要按 grpc_liveness.rs 的风格写一个暂停时钟的测试。测试套件的运行方式见测试。