跳转到内容

一条连接的一生

源码文件:48 个 · 核对版本 Etemenanki 596916d · katana v3.0.1
  • Etemenanki/app/src/serve.rs
  • Etemenanki/app/src/connector.rs
  • Etemenanki/app/src/flow.rs
  • Etemenanki/app/src/router.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/inbound/tun.rs
  • Etemenanki/app/src/transport.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • Etemenanki/app/src/outbound/freedom.rs
  • Etemenanki/app/src/outbound/udp_fanout.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/socks/handshake.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/transports/keepalive.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/connector.rs
  • Etemenanki/protocols/src/hysteria/connection.rs
  • Etemenanki/protocols/src/wireguard/connector.rs
  • Etemenanki/protocols/src/wireguard/device.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/relay.rs
  • Etemenanki/environment/src/routing.rs
  • Etemenanki/environment/src/dial/tcp.rs
  • katana/src/serve.rs
  • katana/src/connector.rs
  • katana/src/manager/proxy.rs
  • katana/src/meter.rs
  • katana/src/router.rs
  • katana/src/outbound/mod.rs

本页跟随一条 TCP 连接穿过 etemenanki-app:从监听器 accept 这个 socket 开始,到最后一个字节排空、任务结束为止。连接途经的每个函数、类型和 effect 都按它遇到的先后顺序列出,每一步都链接到深入介绍该组件的页面。

如果你准备修改服务路径、某个协议核心(core)、connector(连接器)或某个出站,请先读本页:它说明了哪个任务拥有什么、每个超时位于何处,以及协议核心在每个时刻看到的是哪个事件。后面几节依次说明 UDP 有什么不同、两个不接受 TCP socket 的入站(Hysteria 2 和 TUN)有什么不同,以及 katana 有什么不同:katana 运行的是同一个内核,只是换上了自己的 connector。

  1. accept 与准入。 run_stream_inbound accept 这个 socket,并以不等待的方式获取两个信号量许可(permit);拿不到就直接丢弃 socket。
  2. 传输层 accept。 serve_socket 运行 InboundTransport::accept:TLS、WebSocket 升级或 HTTP/2 preface,受 TRANSPORT_HANDSHAKE_TIMEOUT 限制。gRPC 为每个 HTTP/2 stream 产出一条字节流。
  3. 选择协议核心。 serve_connection 构建 AppConnector,并选出该协议的协议核心。SOCKS 则交给它自己的驱动。
  4. 驱动与握手看门狗。 drive 把协议核心包进 ProxyServerRuntime,并在协议核心进入已建立状态之前施加 HANDSHAKE_TIMEOUT。
  5. 解析请求。 协议核心解析并认证请求,然后推入携带 Flow 的 Effect::Open。
  6. 路由。 运行时在应用 Open 的过程中调用 AppConnector::connect,后者把 Flow 及其 FlowContext 转换为 RouteTarget,再向 Router 要一个出站。
  7. 拨号。 出站的 connect_stream future 负责拨号:直连(freedom)、经由 TransportConnector 之上的代理客户端运行时、经由 WireGuard,或经由 Hysteria 2。
  8. 已连接。 运行时投递 Event::Connected;协议核心暂存它的回复,排在 Open 后面的 forward 随之发出。
  9. 中继。 客户端字节变成 Forward effect;从出站读到的数据变成 Event::Outbound,并暂存以发回客户端。
  10. 半关闭与结束。 每个 EOF 都会变成对另一侧的半关闭;两侧都结束后,协议核心推入 Finish,暂存区排空,任务结束,同时丢弃它持有的许可克隆。
sequenceDiagram
  participant C as 客户端
  participant L as accept 循环
  participant T as serve_socket
  participant D as drive
  participant R as ProxyServerRuntime
  participant K as 协议核心
  participant A as AppConnector
  participant X as 出站
  C->>L: TCP 连接
  L->>L: try_acquire 会话许可和握手许可
  L->>T: spawn_scoped(serve_socket)
  T->>C: TLS、WebSocket 升级或 h2 preface
  T->>D: sink(TransportStream), spawn_scoped(serve_connection)
  D->>R: new(stream, core, connector).showing_progress()
  C->>R: 请求字节
  R->>K: Event::Transport(unparsed)
  K-->>R: Effect::Open(key, Flow), Forward(rest)
  Note over D,K: is_established() 为 true,drive 丢弃握手许可
  R->>A: 应用 Open 时调用 connect(flow)
  A->>A: route_target, Router::route
  A-->>R: 装箱的 connect_stream future
  R->>X: 轮询 future:解析、拨号、上游握手
  X-->>R: Outbound::Stream
  R->>K: Event::Connected
  K-->>R: 暂存回复
  R->>X: 排队的 Forward(poll_write, poll_flush)
  R->>C: 暂存的回复
  loop 中继
    C->>R: 字节进入读缓冲区
    R->>K: Event::Transport
    K-->>R: Forward(range)
    R->>X: poll_write
    X->>R: 字节进入 scratch
    R->>K: Event::Outbound
    K-->>R: 暂存
    R->>C: 从暂存区 poll_write
  end
  C->>R: EOF
  R->>K: Event::TransportEof
  K-->>R: Shutdown(key)
  R->>X: poll_shutdown
  X->>R: EOF
  R->>K: Event::OutboundEof
  K-->>R: ShutdownTransport, Finish
  R->>C: 排空暂存区,再关闭写方向
  R-->>D: stream 结束
  D-->>L: 任务结束,许可被丢弃

整条连接只存在于很少几个任务中。既没有按方向划分的任务,客户端侧和出站侧之间也没有 channel:一个 ProxyServerRuntime 拥有客户端传输层、协议核心以及协议核心打开的每一个出站,并亲自在它们之间搬运字节。

任务 由谁 spawn 拥有 何时结束
run_stream_inbound app/src/instance.rs → spawn_generation,每个入站一个 监听器、两个信号量 generation(一代实例)的 CancellationToken 被取消
serve_socket accept 循环,通过 spawn_scoped 已 accept 的 socket,直到传输层产出 stream 为止;对 gRPC 而言是整条 HTTP/2 连接 传输层已产出它的 stream,或 HTTP/2 连接结束
serve_connection serve_socket 的 sink,通过 spawn_scoped(Unix socket 则就地 await) ProxyServerRuntime:传输层、协议核心、出站、三个缓冲区 运行时完成或失败,或 generation 被取消

spawn_scoped 让 future 与 token.cancelled() 竞速,因此取消一个 generation 会丢弃上面所有 future,连同它们拥有的每个 socket。参见 generation 与热重载。

少数出站依赖后台任务:每个拨出的 gRPC 传输层都会 spawn 自己那条 HTTP/2 连接的驱动任务;Hysteria 2 客户端为其所有流共享的 QUIC 连接运行一个 HTTP/3 驱动任务(服务端中继 UDP 时还有一个数据报泵);WireGuard 出站为每条隧道运行一个设备驱动任务。代理客户端运行时(ProxyClientRuntime)本身则作为出站的 AsyncRead 和 AsyncWrite,在连接自己的任务里被轮询。

app/src/serve.rs 中的服务路径:

pub fn spawn_scoped<F>(token: CancellationToken, fut: F) -> JoinHandle<()>
where
F: Future + Send + 'static,
F::Output: Send;
pub async fn run_stream_inbound(
inbound: Arc<StreamInbound>,
tag: CompactString,
listener: StreamListener,
router: Arc<Router>,
token: CancellationToken,
);
async fn serve_socket(
inbound: Arc<StreamInbound>,
socket: AcceptedSocket,
ctx: FlowContext,
router: Arc<Router>,
token: CancellationToken,
session: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
);
async fn serve_connection<S>(
inbound: Arc<StreamInbound>,
stream: S,
local_ip: Option<IpAddr>,
ctx: FlowContext,
router: Arc<Router>,
session: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
) where
S: AsyncRead + AsyncWrite + Unpin + Send + 'static;
async fn drive<const BUF: usize, Core, S>(
stream: S,
core: Core,
connector: AppConnector,
handshake: Arc<OwnedSemaphorePermit>,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
Core: ProxyCoreDecode<Target = Flow, Error = io::Error, TransportAddr = ()> + Established;
pub trait Established {
fn is_established(&self) -> bool;
}

解码后的连接是什么,以及它从哪里进入(protocols/src/flow.rs、app/src/flow.rs):

pub struct Flow<T> {
pub destination: Destination,
pub user: NetworkUser<T>,
pub sniffed: Option<SniffedBehavior>,
pub source: Option<IpAddr>,
}
pub type Flow = etemenanki_protocols::flow::Flow<()>;
pub struct FlowContext {
pub inbound_tag: CompactString,
pub source: Option<IpAddr>,
}

每个流式入站拨号时都经过的 connector(concepts/src/link.rs、app/src/connector.rs):

pub trait Connector<Target> {
type Stream: AsyncRead + AsyncWrite + Unpin;
type Datagram: DatagramLink;
type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;
}
pub struct AppConnector {
pub router: Arc<Router>,
pub ctx: FlowContext,
}
impl Connector<Flow> for AppConnector {
type Stream = OutboundStream;
type Datagram = FanOutLink;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow) -> ConnectFuture;
}

路由与拨号(app/src/router.rs、environment/src/routing.rs、app/src/outbound/mod.rs):

pub fn route_target<'a>(flow: &'a Flow, ctx: &'a FlowContext) -> routing::RouteTarget<'a>;
pub fn route(&self, target: RouteTarget<'_>) -> Arc<Outbound>;
pub fn connect_stream(&self, flow: Flow) -> StreamFuture;
pub fn connect_datagram(&self, flow: Flow) -> DatagramFuture;

运行时,以及它所驱动的协议核心契约(concepts/src/runtime.rs、concepts/src/core.rs):

pub struct ProxyServerRuntime<const BUF_SIZE: usize, Core, Trans, Conn, Mode = ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Conn: Connector<Core::Target>;
pub fn new(transport: T, core: Core, connector: Conn) -> Self;
pub fn showing_progress(
self,
) -> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>;
fn handle(
&mut self,
event: Event<'_, Self>,
effects: &mut Effects<'_, Self>,
) -> Result<usize, Self::Error>;

run_stream_inbound 用 tokio::select! 循环等待 generation 的 token 和 StreamListener::accept。StreamListener 要么是 TcpListener,要么是 Unix 监听器;只有 TCP 监听器会报告对端地址,这个地址成为 FlowContext.source,让 source_cidr 规则有东西可匹配。

准入发生在读取任何字节之前,而且从不等待:

  1. sessions.clone().try_acquire_owned() 从每入站的活动连接信号量(MAX_LIVE_CONNECTIONS_PER_INBOUND,65,536)中取一个名额。
  2. handshakes.clone().try_acquire_owned() 从每入站的握手信号量(MAX_HANDSHAKES_PER_INBOUND,2,048)中取一个名额。

只要有一个信号量已满,socket 就当场被丢弃,已经拿到的许可也一并释放,循环在 debug 级别记录日志(dropping inbound connection; live connection limit reached 或 … handshake limit reached)。选择拒绝而不是等待,可以让循环持续 accept;这样过载的入站会主动丢弃连接,而不是让内核的 listen backlog 悄无声息地被填满。

活动连接许可以 Arc<OwnedSemaphorePermit> 的形式传递。一个 socket 可以产出多条字节流(gRPC 为每个 HTTP/2 stream 产出一条),每条流由各自的任务服务,这些任务每个都持有活动连接许可的一个克隆:只有当服务该 socket 的最后一个任务结束时,名额才会归还给信号量。

accept 错误由 should_backoff_accept_error 分类:ConnectionAborted 和 Interrupted 立即重试;其他错误(例如文件描述符耗尽)在 warn 级别记录,随后休眠 ACCEPT_ERROR_BACKOFF(100 ms),这次休眠本身也可以被取消。token 被取消时循环退出,StreamListener::release 删除 Unix socket 文件。

延伸阅读:服务 和 限制。

serve_socket 把 TCP socket 交给入站的 InboundTransport(protocols/src/transports/accept.rs):

pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);

accept 先在 socket 上启用 TCP keepalive(set_keepalive:TCP_KEEPALIVE_IDLE(120 s)后发出第一个探测,之后每隔 TCP_KEEPALIVE_INTERVAL(30 s)一次,TCP_KEEPALIVE_RETRIES(3)次后放弃),然后在 within 中运行传输层自身的握手;within 施加 TRANSPORT_HANDSHAKE_TIMEOUT(10 s),超时则以 tls handshake timed out(或 websocket …、grpc …)失败:

InboundTransport 10 s 限制内的握手 交给 sink 的 stream
Tcp 无 恰好一条,立即交出
Tls(ServerConfig) TLS accept 恰好一条
Ws { route, tls } 可选的 TLS,然后是 WebSocket 升级(检查路径和 Host) 恰好一条
Grpc { paths, tls } 可选的 TLS,然后是 HTTP/2 服务端握手 路径已知的每个已 accept 的 HTTP/2 stream 各一条,直到连接结束

对 gRPC 而言,serve_h2 在握手之后继续驱动这条 HTTP/2 连接。每个已 accept 的 stream,只要其路径能被 GrpcPaths::classify 识别(/{service}/Tun 或 /{service}/TunMulti),就会收到 gRPC 响应头,并作为 TransportStream::Grpc 交给 sink;其他路径一律以 REFUSED_STREAM 重置。Liveness 监视这条连接,对已经空闲或不再响应 ping 的对端放弃连接。

serve_socket 中的 sink 为每条产出的 stream spawn 一个 serve_connection,使用同一个 token,带上该 socket 的许可及其 local_addr 的 IP(只有 SOCKS 驱动会用到)。Unix socket 完全跳过传输层:配置构建器拒绝在 Unix 监听器上设置任何 stream network 或 security(inbound <tag>: protocol http over a unix socket does not support stream network "ws"),因此 serve_socket 直接 await serve_connection,local_ip 和 source 都为 None。

传输层如果在产出任何 stream 之前就失败,serve_socket 会记一条 debug 日志(inbound transport failed: …)后结束;客户端根本到不了协议核心。

延伸阅读:TCP 与 TLS 传输层 和 WebSocket 与 gRPC 传输层。

serve_connection 首先构建 connector,因为每个协议核心都要通过它拨号:

let connector = AppConnector {
router,
ctx: ctx.clone(),
};

然后它对 StreamInbound.protocol(每个 generation 构建一次的 StreamProtocol)进行 match,为这条连接构造一个全新的协议核心。协议核心通过 Arc 借用入站共享的用户表,运行时的缓冲区大小就是协议核心自己的 BUF_SIZE 常量:

StreamProtocol 驱动 协议核心构造函数 BUF_SIZE
Socks(SocksInbound<()>) SocksInbound::serve 无:专用驱动 RELAY_BUF,BidirectionalConnection 每个方向 16 KiB
Http(..) drive HttpCore::new(config, sniff, source) MAX_HEAD,64 KiB
Trojan(..) drive TrojanCore::new(validator, sniff, source) 16 KiB
Vless(..) drive VlessCore::new(validator, sniff, source) 16 KiB
Vmess(..) drive VMessCore::new(validator, now_unix, sniff, source) 32 KiB
Shadowsocks(..) drive ShadowsocksCore::new(resolved, sniff, source) 20 KiB
Ss2022 { config, validator } drive Ss2022Core::with_system_clock(config, validator, sniff, source) 32 KiB

时钟是构造函数参数(VMess 用 now_unix,Shadowsocks 2022 用系统时钟),而不是由运行时注入的东西,因此协议核心在测试中保持确定性。

SOCKS 是唯一不适配 ProxyCoreDecode 的服务端:它的 UDP 部分位于第二个 socket 上,控制连接只负责让它保持存活,而且这个 socket 只接收控制连接那个客户端的数据报(见 UDP 有何不同)。SocksInbound::serve 在自己的 HANDSHAKE_TIMEOUT 下运行 SOCKS 4/4a/5 握手,用同一个 AppConnector 调用 connector.connect(flow),只在连接成功后才写回授权回复(失败则写回拒绝码;嗅探 IP 目标是例外,见第 5 步),并用受 RELAY_IDLE_TIMEOUT 保护的 BidirectionalConnection 进行中继。从第 4 步开始描述的都是协议核心路径。

延伸阅读:服务;SOCKS 驱动见 SOCKS。

drive 以 stream 模式构建运行时:

let mut runtime =
ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress();

showing_progress 把运行时从 Future 变成一个 Traffic 增量的 Stream,每完成一个工作单元产出一项。app 并不统计流量;它使用 stream 模式,是为了能在步与步之间检查协议核心。

客户端发送第一个字节之前,运行时不会投递任何事件,而协议核心自己的 Timing 只在第一个字节事件到来时才启动握手截止时间。如果客户端连上后一言不发,本来永远不会超时,因此 drive 充当看门狗:只要 runtime.core().is_established() 为 false,每次 runtime.next() 都包在 tokio::time::timeout(HANDSHAKE_TIMEOUT, …) 里。超时后 drive 返回 inbound handshake timed out after 10s,连接被丢弃。

两个截止时间互相补位。drive 负责抓住从不说话的客户端;第一个字节到达后,Timing::touch 只启动一次 HANDSHAKE_TIMEOUT(阶段仍为 Handshake 时不会重新启动),因此一个逐字节慢慢发送请求的客户端,仍会在第一个字节之后 10 s 撞上协议核心的截止时间,协议核心以 client did not complete its request in time 失败。

is_established() 来自 Timing::is_established:协议核心处于 Phase::Relay 或 Phase::Closing 时为 true。协议核心在接受请求时进入 Relay:对普通 TCP 请求,是在推入 Open 的同一次调用中,早于拨号完成;对 UDP 或 mux 请求,是在头部认证通过后立即进入;对需要嗅探的请求,则要等嗅探结束。此时 drive 丢弃它持有的握手许可克隆(handshake.take()),然后一直轮询运行时直到结束;从这里开始,由协议核心自己的截止时间管理连接。

延伸阅读:服务端运行时 和 协议基础。

传输层读到的数据落入运行时的读缓冲区(up,一个 ReadBuffer<BUF_SIZE>)。运行时把整个未解析区域作为 Event::Transport(&mut [u8]) 交给协议核心;协议核心返回它消费了开头多少字节。需要更多数据的协议核心返回 Ok(0),运行时保留未消费的尾部并继续读取。如果未解析区域已占满整个缓冲区而协议核心仍一个字节都不消费,运行时以 RuntimeError::FrameTooLarge 失败。

TrojanCore 是最朴素的例子。在第一个字节事件上,Timing::touch 推入 Effect::SetDeadline(Some(HANDSHAKE_TIMEOUT))。一旦 parse_request_header 拿到完整的头部,协议核心就在它的 Validator 中查找密码哈希,构建流并打开它:

let flow = Flow::new(header.destination, user, self.source);
// ...
fx.push(Effect::Open {
key: FlowKey::Direct,
target: flow,
});

open_tcp 同时进入 Phase::Relay,这会把截止时间重新设为 RELAY_IDLE_TIMEOUT,然后这次调用返回头部长度。运行时立刻把缓冲区剩余部分再次交给协议核心,协议核心为随头部一起到达的载荷推入 Forward { key: FlowKey::Direct, range }。这个 forward 需要等待:运行时严格按顺序应用 effect,而指向仍在连接中的 key 的 forward 会让队列停下。

未知用户会让协议核心返回错误(trojan: invalid user,PermissionDenied),运行时将其报告为 RuntimeError::Core;连接在没有任何回复的情况下结束。

当入站开启了 sniff,且请求指定的是 IP(worth_sniffing)时,协议核心不会立即打开。它进入 Phase::Sniff(截止时间为 SNIFF_TIMEOUT,300 ms),把最初的载荷字节复制到 SniffPrefix 中,最多 SNIFF_LIMIT(4 KiB),直到出现 TLS SNI 或 HTTP Host、预算用完,或截止时间到期。随后它设置 flow.sniffed,推入 Open,并用 Effect::ForwardHeld 从其 held 缓冲区转发收集到的字节。嗅探得到的域名只提供给路由器;拨号仍然发往 flow.destination。

对于客户端要先等到回复才发送载荷的协议(HTTP CONNECT、SOCKS、Hysteria 2),在嗅探 IP 目标时会先回复“已连接”,再收集前缀。

延伸阅读:服务端协议核心、嗅探,以及各协议自己的页面,例如 Trojan 或 VLESS。

路由在运行时的任务中同步执行,发生在运行时应用 Effect::Open 的那一刻。在 drive_effects 中,运行时取出这个 effect,拒绝已经存活的 key(RuntimeError::DuplicateKey),为该 key 创建 KeyWaker,调用 self.connector.connect(target),把返回的 future 装箱存入 LinkState::Connecting,然后调度该 key。

AppConnector::connect 按流的网络类型分流:

fn connect(&mut self, flow: Flow) -> ConnectFuture {
if flow.destination.network == DialNetwork::Udp {
let link = FanOutLink::new(self.router.clone(), self.ctx.clone(), flow);
return Box::pin(async move { Ok(Outbound::Datagram(link)) });
}
let outbound = self.router.route(route_target(&flow, &self.ctx));
let fut = outbound.connect_stream(flow);
Box::pin(async move { fut.await.map(Outbound::Stream) })
}

route_target 把流适配为路由模型的 RouteTarget:

RouteTarget 字段 取自 使用者
remote、port flow.destination 域名、cidr、geoip 和 port 匹配条件
sniffed_domain flow.sniffed 每个域名匹配条件(包括 geosite)除了 remote 之外也会尝试它
inbound_tag ctx.inbound_tag inbound_tag
source flow.source.or(ctx.source) source_cidr
network flow.destination.network(Tcp 或 Udp) network

流自己的来源优先于监听器的来源。对 TCP 入站而言两者是同一个地址;QUIC 或 TUN 入站服务许多对端,只有流本身知道它来自哪一个。

Router::route 调用 RouteTable::pick:按顺序尝试规则,第一条有任一匹配条件命中的规则胜出,否则返回默认出站。结果是一个在整个 generation 内共享的 Arc<Outbound>。

延伸阅读:路由模型。

Outbound::connect_stream(flow) 返回一个装箱的 StreamFuture;在运行时于 service_key 中轮询它之前,不会发生任何拨号。Outbound::Balanced 在这里调用 Balancer::select(),因此宕掉的成员从下一个流开始就会被跳过。Outbound::Blackhole 立即解析为 OutboundStream::Blackhole,读取时返回 EOF,写入则被吞掉。

FreedomConnector::dial 用 destination_to_socketaddrs(使用构建该出站时指定的 Resolver 和 AddressFamilyStrategy)解析 flow.destination,然后 TcpDialer::connect_any 逐个尝试这些地址,每个地址受 DEFAULT_CONNECT_TIMEOUT(10 s)限制,并在一个错误中返回所有失败(failed to connect to any address (…))。结果是 OutboundStream::Tcp。

延伸阅读:出站、客户端运行时、拨号器、WireGuard、Hysteria 2 客户端。

8. 已连接,以及协议核心暂存的回复

Section titled “8. 已连接,以及协议核心暂存的回复”

每当 key 的 waker 被唤醒,service_key 就轮询连接 future。在 Ready(Ok(outbound)) 时,槽位变为 LinkState::Stream,key 被重新调度以便尝试第一次读取,协议核心收到 Event::Connected { key }。协议核心如何处理它因协议而异:

协议核心 收到 Connected(TCP)时 收到 ConnectFailed 时
HttpCore(CONNECT) 暂存 HTTP/1.1 200 Connection established 暂存 HTTP/1.1 502 Bad Gateway(除非因嗅探 IP 目标已经发出了 200),然后 ShutdownTransport 和 Finish
VlessCore 暂存 2 字节的 RESPONSE_HEADER 不回复,直接关闭
VMessCore 暂存加密封装后的响应头 不回复,直接关闭
TrojanCore 什么都不做:Trojan 没有回复 关闭
ShadowsocksCore 什么都不做 关闭
Ss2022Core 暂时什么都不做:响应头会随第一批下行字节一起加密封装;若没有下行字节,则在出站 EOF 时封装 关闭

协议核心返回后,handle_meta 立即运行 drive_effects,排在 Open 后面等待的 forward 被写入出站。暂存的回复在下一次 poll_transport_write 时发给客户端。

在 Ready(Err(error)) 时,运行时调用 forget_key,它移除该槽位并丢弃所有仍指向该 key 的排队 effect,然后投递 Event::ConnectFailed { key, error }。连接失败是交给协议核心应对的事件,而不是运行时错误。单 stream 的协议核心通过 Passthrough::on_outbound_gone 应对,它推入 ShutdownTransport 和 Finish。

UDP 和 mux 请求会更早得到回复:VLESS 和 VMess 在请求认证通过后就立即暂存响应头,因为这些流之后才会按数据包或按子会话打开。

延伸阅读:VLESS、VMess 线格式、HTTP、Shadowsocks 2022。

从这里开始,运行时在两个方向之间交替工作。每次 step 按以下顺序处理:截止时间、仍在等待某个出站的 effect、发往客户端的暂存写入,然后是客户端读取和已就绪的出站,两者轮流进行(prefer_transport),谁也不会饿死对方。

上行(客户端到目标)。 一次传输层读取填充 up;Event::Transport 把未解析的字节交给协议核心;处于中继状态的协议核心推入 Forward { key, range }(对明文协议核心,Passthrough::on_transport 会转发整段切片)。排队中的 forward 会钉住传输层缓冲区:在这个 forward 写出之前,运行时不会再从客户端读取。drive_effects 先用 poll_write 写入,再 poll_flush;如果出站返回 Pending,队列就停在那里,客户端也不会被读取。背压就是这样传递到客户端的。

下行(目标到客户端)。 出站只有在其 KeyWaker 被唤醒后才会得到处理,而且只有当暂存缓冲区还剩 STAGING_RESERVE + 1 字节空间时才会被读取。读取的数据进入 scratch(最多为预留量之上的剩余空间),并以 Event::Outbound { key, data } 的形式到达协议核心,协议核心必须整块消费:它暂存这些字节,对加密协议则先加密封装。poll_transport_write 把暂存区写给客户端。不再读取的客户端会把暂存区填满,运行时随之停止读取出站。

空闲。 Timing::touch 在每个字节事件开始时运行,并重新启动 RELAY_IDLE_TIMEOUT(300 s)。当截止时间到期而两个方向都没有任何数据流动时,Timing::expired 返回 Expired::Idle 并推入 Finish。

当出站是代理客户端运行时时,协议核心的 Forward 落入 ProxyClientRuntime::poll_write,后者用 codec 把字节加密封装到它自己的暂存缓冲区,再写给上游;poll_read 就地解开帧并把明文复制出来。这一切都在连接自己的任务里运行。

延伸阅读:服务端运行时 和 链路与类型。

单 stream 的协议核心用一个 Passthrough 来记录关闭状态:

事件 Passthrough 调用 推入的 effect
Event::TransportEof(客户端发送完毕) on_transport_eof Shutdown { key };如果出站已经结束,再推入 Finish
Event::OutboundEof { key }(目标发送完毕) on_outbound_eof ShutdownTransport;如果客户端已经结束,再推入 Finish
Event::ConnectFailed 或 Event::OutboundError on_outbound_gone ShutdownTransport、Finish

effect 按顺序应用,因此 Shutdown { key } 只有在它之前排队的所有 forward 都写出后才会到达出站(poll_shutdown)。ShutdownTransport 只有在暂存区排空并 flush 之后才关闭客户端的写方向。读写两个方向都已关闭的槽位会从运行时的 outbounds 映射中移除,这会丢弃(并关闭)该出站。

Finish 之后,运行时不再从任何一侧读取,只排空暂存区。当协议核心已结束、暂存区为空、且请求过的传输层关闭已完成时,step 返回 None。进度 stream 随之结束,drive 返回 Ok(()),serve_connection 返回,运行时及其拥有的一切被丢弃,连接持有的会话许可克隆也一并释放;当该 socket 的最后一个克隆消失时,许可归还给信号量。

失败时的结束方式相同,只是更早。drive 把任何 RuntimeError 转换为 io::Error,serve_connection 在 debug 级别记录它,附带协议的 StreamProtocol::name() 和来源的 Debug 形式,例如 trojan connection from Some(192.0.2.10) ended: proxy core: trojan: invalid user。

延伸阅读:Passthrough 和关闭类 effect 见 服务端协议核心,排空过程见 服务端运行时。

UDP 流以同样的方式到达 connector,即一个 Flow 为 DialNetwork::Udp 的 Effect::Open,但 AppConnector 既不为它选路也不为它拨号。它立即返回一个 FanOutLink,因此协议核心在下一次轮询时就会收到 Connected,之后每个数据包都单独选路:

flowchart TB
  core["协议核心:Effect::SendTo(key, to, range)"]
  send["FanOutLink::poll_send_to"]
  rt["route_target(flow.toward(to), ctx)"]
  route["Router::route:Arc of Outbound"]
  known{"该出站的子链路已打开?"}
  write["子链路 poll_send_to,移到最近使用"]
  busy{"另一条子链路正在打开?"}
  opendg["Outbound::connect_datagram(flow.toward(to))"]
  wait["Pending:运行时保留排队的 SendTo"]
  drop["打开失败:丢弃数据包"]
  core --> send --> rt --> route --> known
  known -->|是| write
  known -->|否| busy
  busy -->|是| wait
  busy -->|否| opendg
  opendg -->|成功| write
  opendg -->|出错| drop
  • 标识。 子链路以所选出站的 Arc::as_ptr 为键。路由器的表在整个 generation 期间持有这些 Arc,因此指针就代表“同一个出站”。
  • 路由输入。 每个数据包都以 self.flow.toward(to) 选路:沿用该关联的用户和来源、使用数据包自己的目标(其网络为 Udp),且不带嗅探域名(toward 会将其重置)。
  • 一次只打开一条。 同一时间只打开一条子链路。发往另一个出站的数据包会返回 Pending,直到那次打开完成,期间运行时保留它排队的 SendTo。打开失败时丢弃数据包(Ok(buf.len())),这样 UDP 关联不会因一个不可达的对端而卡住。
  • 上限。 最多同时打开 MAX_SUBS(64)条子链路;打开第 65 条时,淘汰最久未发送过的那条。
  • 回复。 poll_recv_from 从上一次接收停下的位置开始轮转轮询各子链路,返回第一个数据包及其来源,该数据包以 Event::Datagram { key, from, data } 的形式到达协议核心。接收失败的子链路会被移除;关联继续使用其余的子链路。

Outbound::connect_datagram 打开针对单个出站的链路:freedom 用一个双栈 socket;Trojan、VLESS 和 VMess 用 ProxyClientRuntime 之上的协议数据报 codec;SOCKS 用 SocksUdpLink(只有当发送方的 endpoint 等于服务端给出的中继的 endpoint 时,它才把该数据报当作回复,因此双栈 socket 也能收到 IPv4 中继的回复);WireGuard 用一个隧道关联;Hysteria 2 用共享 QUIC 连接上的一个会话。HTTP、Shadowsocks 和 Shadowsocks 2022 出站返回 Unsupported(… carries no datagrams),因此发往它们的数据包会被丢弃。

SocksInbound::serve 通过 handshake_with_udp_source 运行握手,它除了返回 Handshake,还返回 UDP ASSOCIATE 请求中的 DST.ADDR 和 DST.PORT(其他请求为 None)。对于 UDP ASSOCIATE,associate 在绑定中继 socket(hub)之前,用控制连接的 source、请求声明的这个来源以及中继 IP 构建一个 ExpectedSender:

  • 只认一个客户端。 经 TCP 时,每个数据报都必须来自控制连接的 IP,比较时由 endpoint(protocols/src/socks/protocol.rs)转成规范形式,因此 IPv4 映射的 IPv6 地址等同于对应的 IPv4 地址。来自其他发送方的数据报在解析之前就被丢弃,也不会重置该关联的 RELAY_IDLE_TIMEOUT。
  • 固定端口。 第一个被转发的数据报(能解析且带有载荷)会固定端口,此后回复只发往这个地址。如果请求给出的是控制连接自己的 IP 和一个非零端口,则一开始就固定该端口。请求给出其他任何来源时,该来源被搁置,既不信任也不拒绝。
  • 没有可认定的客户端。 经 Unix socket 时 source 为 None,因此请求必须给出确切的 IP 地址和非零端口。不满足这一点的请求,以及 hub 所在地址族收不到客户端数据报的情形(hears:IPv4 或 IPv4 映射的 hub 只收 IPv4,未指定的 IPv6 hub 两者都收,其他 IPv6 hub 只收 IPv6),都会在绑定任何 socket 之前收到 SOCKS5 回复 STATUS_NOT_ALLOWED(0x02)。

驱动在转发第一个数据报时,以控制连接的 source 作为流的来源,向同一个 AppConnector 请求一个 UDP 流,得到同样的 FanOutLink,并在自己的 select 循环中驱动它。

延伸阅读:出站、路由模型、SOCKS:关联的客户端;mux 承载中的 UDP 见 mux.cool 与 XUDP。

这两个入站都不接受 TCP socket,因此第 1 到第 4 步被入站自己的运行循环取代。两个循环都接收一个闭包,它根据客户端地址为一个运行时构建 AppConnector(循环每启动一个运行时就调用它一次);从 Effect::Open 开始,上面的一切原样适用:

let make = move |ip| AppConnector {
router: router.clone(),
ctx: FlowContext {
inbound_tag: tag.clone(),
source: Some(ip),
},
};
Hysteria 2(run_hysteria_inbound → Hy2Inbound::run) TUN(run_tun_inbound → TunInbound::run)
接受什么 一个 UDP socket 上的 QUIC 连接,每条连接占用一个 max_connections 许可(监听器满时拒绝握手) 由设备上的用户态 IP 栈终结的 TCP 连接和 UDP 流
认证 每条 QUIC 连接一次,通过 HTTP/3 无;TunConfig.user 用于归属
TCP 每个已分类的代理 stream 一个基于 Hy2StreamCore 的 ProxyServerRuntime,各占一个 circuit_permits 许可 每条 TCP 连接一个基于 PassthroughCore 的 ProxyServerRuntime,占用一个 max_flows 许可
UDP 监听器中继 UDP 时,每条已认证的 QUIC 连接一个基于 Hy2UdpCore 和 QuicDatagrams(over_datagrams)的运行时 启用 UDP 时,每个客户端来源一个基于 TunUdpCore 和 TunUdpLink 的运行时,占用一个 max_flows 许可,该来源的新流通过 FLOW_QUEUE(16)channel 送入
取消时 运行循环返回,然后 Hy2Inbound::shutdown 等待 UDP 端口释放 运行循环返回,然后 TunInbound::shutdown 等待设备关闭

PassthroughCore 在第一个字节到来之前就知道目标,因为 IP 栈已经知道客户端要去哪里。它在第一个事件上打开,原样中继,可选嗅探。由于客户端说话之前运行时不会投递任何事件,TUN 的 serve_stream 施加了与 drive 相同的静默客户端看门狗:在协议核心建立之前,每一步都包在 HANDSHAKE_TIMEOUT 中(tun: the client never spoke)。

延伸阅读:Hysteria 2 服务端、TUN。

katana 使用同一组内核 crate,整体形态也相同:accept、传输层、每条 stream 一个运行时、一个负责路由和拨号的 Connector。差异在于准入、运行时外围的看门狗,最重要的是 connector。

etemenanki-app(app/src/serve.rs) katana(src/serve.rs)
任务作用域 spawn_scoped(token, fut) 返回一个 JoinHandle Scope 把 token 与一个 TaskTracker 配对;Scope::shutdown 等待所有任务
活动连接 每个监听器 MAX_LIVE_CONNECTIONS_PER_INBOUND(65,536),try_acquire 每个监听器 MAX_LIVE_CONNECTIONS_PER_NODE(65,536),try_acquire;拒绝计为一次握手失败
stream 上的协议 SOCKS、HTTP、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022 VMess、VLESS、Trojan、Shadowsocks、Shadowsocks 2022,每个协议核心都用 stream 到达时的当前用户表构建
握手看门狗 建立之前,每次 runtime.next() 外包一个 HANDSHAKE_TIMEOUT 建立之前,整个循环外包一个 HANDSHAKE_TIMEOUT;届时丢弃 pre-auth 许可
握手之后 轮询运行时直到结束;由协议核心自己的截止时间管理 对运行时、用户下线(until_retired)和 PROGRESS_WATCHDOG(RELAY_IDLE_TIMEOUT + 60 s = 360 s)做 select!;每当 Traffic 增量非零就重置看门狗

katana 的 accept 循环还对处于传输层握手中的 socket 设了上限,即 MAX_TRANSPORT_STAGES_PER_NODE(2,048)个名额;名额不足时它会等待,而不是拒绝 socket。对传输层产出的每条 stream,sink 调用 ProxyManager::accept_stream:它以不等待的方式从节点的 pre-auth 信号量(MAX_PREAUTH_STREAMS_PER_NODE,512)取一个名额,拒绝计为一次握手失败,然后带着该 socket 活动连接许可的一个克隆 spawn serve_stream。协议核心建立后,drive 丢弃 pre-auth 许可。每条 stream 的协议核心都用 stream 到达时的当前用户表(protocol.load_full())构建,因此用户刷新永远不会在握手进行中替换掉它所用的表。

katana 打开的每个流,无论是 TCP 请求、mux 子流、UDP 关联还是 Hysteria 2 stream,都要经过 KatanaConnector::connect:

impl Connector<Flow<UserTag>> for KatanaConnector {
type Stream = Metered<OutboundStream>;
type Datagram = FanOut;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;
}
flowchart TB
  flow["Flow of UserTag"]
  admit{"Admission::admit(tag)"}
  lease["把用户的租约发布到连接的 LeaseSlot"]
  gate["Gate::new(counter, lease)"]
  udp{"网络是 UDP?"}
  fanout["FanOut:逐包路由、审计和计费"]
  route["Router::route(route_target(dest, sniffed, source))"]
  check{"Outbound::Block 或命中审计?"}
  dial["connect_stream(dest)"]
  metered["Metered::new(stream, gate)"]
  refused["refused():PermissionDenied"]
  flow --> admit
  admit -->|未知用户| refused
  admit -->|已准入| lease --> gate --> udp
  udp -->|是| fanout
  udp -->|否| route --> check
  check -->|是| refused
  check -->|否| dial --> metered
  1. 准入。 Admission::admit 在节点的注册表中查找流的 UserTag,返回该用户的 UserCounter 和租约(一个 CancellationToken);对于面板已不再列出的用户,或现已绑定到另一个 uid 的凭据,返回 None。查找与用户刷新提交时持有的是同一把锁。
  2. 租约。 连接上第一个被准入的流把它的租约发布到连接的 LeaseSlot(watch::Sender<Option<CancellationToken>>),这样用户下线时 drive 就能结束整条连接。
  3. 路由。 katana 自己的 route_target(dest, sniffed, source)(src/router.rs)根据目标(包括其网络类型)、嗅探域名和 flow.source.or(self.source) 构建 RouteTarget。它不设置入站 tag。
  4. 审计。 路由到 Outbound::Block 的流,或 RuleManager::detect 报告为禁止访问(并记到该 uid 名下)的目标,以 PermissionDenied 拒绝。
  5. 拨号与计量。 出站拨号 flow.destination,得到的 stream 包装进 Metered。每次读写都先等待 Gate::poll_open:用户一旦下线,它以 ConnectionAborted 失败;用户的令牌桶处于欠额时,它就休眠。之后再用 Gate::sent 或 Gate::received 为这次传输计费。

UDP 交给 FanOut,它是 katana 对应 FanOutLink 的组件,对每个数据包做路由、审计和计费:被阻断或禁止的数据包直接丢弃且不计费,每个发出的数据包和每个回复都要经过同一个 Gate。

对于 Hysteria 2,run_hysteria 构建的每个 KatanaConnector 都没有租约槽;下线用户的 Hysteria 流之所以会结束,是因为它们的 Metered 出站拒绝再搬运数据。

延伸阅读:katana 的 服务、准入、计量 和 Connector 与 UDP。

不变量 机制 由谁固定
socket 只有在有空闲活动连接名额时才被准入,且只要还有任何服务该 socket 的任务存活,名额就一直被占用 spawn 之前 sessions.try_acquire_owned();Arc<OwnedSemaphorePermit> 被带入每个 serve_connection 无专门测试
generation 的每个任务都随它一起结束 spawn_scoped 对 token.cancelled() 做 select;丢弃运行时即丢弃传输层和每个出站 流式入站无专门测试;app/tests/integration/e2e_hysteria_inbound.rs → a_reload_rebinds_the_udp_port 覆盖了自己拥有 socket 的监听器
一条 gRPC 连接为每个 HTTP/2 stream 产出一条 stream,空闲连接会被放弃 serve_h2 为每个已 accept 的 stream 调用 sink;Liveness protocols/tests/pipeline/transports.rs → one_grpc_connection_carries_many_streams;protocols/tests/unit/transports/accept.rs → a_connection_that_opens_no_stream_is_given_up_on
握手截止时间只启动一次,之后空闲截止时间在每个字节事件上重新启动 Timing::touch、Timing::enter protocols/tests/unit/core/mod.rs → timing_arms_handshake_once_then_idle_per_byte_event、timing_reports_handshake_and_sniff_expiry_to_the_core
与 Open 同一批转发的字节会等待连接完成,held 字节会一直保留到被应用 drive_effects 在 Connecting 槽位处停下;held 区间在应用时才解析 concepts/tests/runtime.rs → held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites
大于读缓冲区的帧会使连接失败,而不是扩大缓冲区 RuntimeError::FrameTooLarge concepts/tests/runtime.rs → frame_larger_than_the_buffer_is_an_error
流自己的来源优先于监听器的来源 route_target 中的 flow.source.or(ctx.source) app/tests/unit/router.rs → the_circuits_own_source_wins_over_the_listeners、without_one_the_listeners_source_is_used、a_circuit_source_works_with_no_listener_source
路由上下文能到达匹配条件 route_target 设置 inbound_tag、source 和 network app/tests/integration/e2e_route_context.rs → inbound_tag_selects_the_route、source_cidr_matches_the_client_address、network_separates_tcp_from_udp
嗅探到的域名可以为以 IP 寻址的流选路 SniffPrefix 设置 flow.sniffed;RouteTarget::sniffed_domain app/tests/integration/e2e_sniff.rs → an_http_host_routes_an_ip_addressed_flow、a_tls_sni_routes_an_ip_addressed_flow、turning_sniffing_off_stops_the_domain_rule_matching
Connected 意味着上游确实接受了该流 ProxyClientConnecting 等待 poll_connected(拨号、握手、flush) concepts/tests/client.rs → upstream_dial_failure_is_connect_failed_not_connected、refused_upstream_handshake_is_connect_failed_not_connected
只有目标建立后才暂存回复;例外是为“客户端需先等回复”的协议嗅探 IP 目标时 回复在 Event::Connected 分支中暂存 protocols/tests/unit/vless/core.rs → tcp_request_opens_and_replies_only_once_connected、a_refused_connect_closes_without_a_reply;protocols/tests/unit/http/core.rs → connect_to_a_domain_answers_200_once_connected、connect_refused_answers_502_and_closes、connect_to_an_ip_with_sniffing_replies_early_and_holds_the_prefix;protocols/tests/unit/vmess/core.rs → tcp_request_opens_and_replies_once_connected
连接失败是事件,而不是运行时错误 先 forget_key,再 Event::ConnectFailed concepts/tests/runtime.rs → connect_failure_reaches_the_core_as_an_event;protocols/tests/unit/trojan/core.rs → connect_failure_ends_the_connection_without_a_reply
卡住的出站会让上行停下,而不是把上行数据缓存起来 排队的 forward 钉住传输层缓冲区 concepts/tests/runtime.rs → stalled_outbound_holds_uplink_but_not_other_downlink
代理出站在连接自己的任务中运行 ProxyClientRuntime 就是出站的 AsyncRead/AsyncWrite concepts/tests/client.rs → server_runtime_relays_through_a_client_runtime_in_one_task
每个 EOF 都会半关闭另一侧,两侧都结束后连接完成 Passthrough、Effect::Shutdown、Effect::ShutdownTransport concepts/tests/runtime.rs → half_close_propagates_both_ways_and_finishes;protocols/tests/unit/core/mod.rs → passthrough_half_closes_each_side_and_finishes_on_the_second
UDP 逐包选路,回复汇合回同一个关联 FanOutLink app/tests/integration/e2e_udp_route.rs → one_association_routes_each_peer_separately、replies_from_several_peers_merge_back_correctly
SOCKS UDP 关联只接收其控制连接那个客户端的数据报,请求也无法放宽这一点 ExpectedSender::admits 和 pin;两端都用 endpoint protocols/tests/unit/socks/server.rs → only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port、an_ipv4_mapped_address_is_the_ipv4_one、a_request_naming_the_peer_pins_its_port_up_front、a_request_naming_any_other_source_is_set_aside、over_a_unix_socket_the_request_must_name_the_exact_source、a_relay_that_cannot_hear_the_client_is_refused;protocols/tests/pipeline/socks.rs → udp_association_ignores_another_ip、udp_association_ignores_another_port_once_pinned、udp_association_holds_to_the_port_the_request_names、udp_association_sets_aside_a_source_it_cannot_hold_to、udp_association_is_not_widened_by_the_request、udp_association_over_a_unix_socket_needs_its_exact_source、udp_association_refuses_a_relay_that_cannot_hear_the_client、udp_link_ignores_datagrams_not_from_the_relay、udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay;protocols/tests/unit/socks/protocol.rs → endpoint_sees_through_ipv4_mapping_and_ignores_flow_info
katana 对每个流都做准入、审计和计量 KatanaConnector::connect、FanOut、Metered katana tests/unit/connector.rs → a_user_the_registry_does_not_know_is_refused、a_blocked_destination_is_refused、a_forbidden_destination_is_refused_and_recorded、an_admitted_stream_is_billed_to_its_user、udp_is_billed_after_routing_and_blocked_packets_are_free
katana 会结束静默、卡住和已下线用户的连接 drive 的握手超时、PROGRESS_WATCHDOG、until_retired katana tests/unit/serve.rs → a_silent_client_is_dropped_at_the_handshake_deadline、a_connection_that_stops_moving_is_dropped、a_retired_users_connection_ends
位置 发生什么 客户端看到什么
accept 时信号量已满 丢弃 socket,记 debug 日志 连接在任何字节之前关闭
accept 错误 退避 100 ms,除非是 ConnectionAborted 或 Interrupted 无影响;继续下一次 accept
传输层握手失败或超过 10 s serve_socket 记录 inbound transport failed TLS 或升级失败,或连接关闭
未知路径上的 gRPC stream REFUSED_STREAM 该 stream 被重置;连接保持
静默客户端 drive 在 HANDSHAKE_TIMEOUT 后超时 关闭
请求格式错误或用户未知 协议核心返回错误;RuntimeError::Core 关闭,没有协议回复
路由 不会失败:总有默认出站兜底 —
SOCKS UDP ASSOCIATE 没有中继可认定的客户端(经 Unix socket 的请求未给出确切地址和端口,或 udp_bind 所在地址族收不到客户端) 在绑定 hub 之前 ExpectedSender::new 失败;PermissionDenied SOCKS5 回复 0x02,然后关闭
SOCKS 中继收到来自关联客户端以外任何人的数据报 在解析之前丢弃,不重置空闲计时器;关联继续 无
拨号或上游握手失败 Event::ConnectFailed;由协议核心应对 HTTP 502、SOCKS 拒绝码、Hysteria 拒绝,或关闭
中继期间出站读写出错 fail_outbound,然后 Event::OutboundError ShutdownTransport 和 Finish:发完已暂存字节后正常关闭
传输层读写出错 RuntimeError::Transport 结束运行时 socket 已经不在了
空闲达到 RELAY_IDLE_TIMEOUT Expired::Idle 推入 Finish 已暂存字节排空,然后关闭
generation 被取消(重载或关闭) 每个作用域内的 future 都被丢弃 旧 generation 上的所有连接被突然关闭

整条路径都通过 drop 实现取消。路径中没有任何地方需要显式调用关闭:丢弃运行时就会关闭传输层和每个出站,丢弃 Arc 许可句柄就会把名额归还给各自的信号量。

常量 值 定义位置 作用于
MAX_LIVE_CONNECTIONS_PER_INBOUND 65,536 app/src/serve.rs 每个流式入站的活动 socket
MAX_HANDSHAKES_PER_INBOUND 2,048 app/src/serve.rs 每个流式入站处于握手阶段的连接
ACCEPT_ERROR_BACKOFF 100 ms app/src/serve.rs 非瞬时 accept 错误后的休眠
TRANSPORT_HANDSHAKE_TIMEOUT 10 s protocols/src/transports/accept.rs TLS、WebSocket 升级、HTTP/2 preface
TCP_KEEPALIVE_IDLE、…_INTERVAL、…_RETRIES 120 s、30 s、3 protocols/src/transports/keepalive.rs InboundTransport::accept 收到的每个 socket 以及 TransportConnector::dial 打开的每个 socket(freedom 的 socket 保持平台默认值)
HANDSHAKE_TIMEOUT 10 s protocols/src/core/mod.rs drive 的看门狗以及协议核心的握手截止时间
SNIFF_TIMEOUT、SNIFF_LIMIT 300 ms、4 KiB protocols/src/sniff/mod.rs 嗅探窗口
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs 中继空闲
DEFAULT_CONNECT_TIMEOUT 每个地址 10 s environment/src/dial/tcp.rs freedom 和 TransportConnector 的 TCP 连接
OPEN_STREAM_TIMEOUT 5 s protocols/src/hysteria/connection.rs 打开一条 Hysteria 2 代理 stream
WORK_BUDGET 64 个单元 concepts/src/runtime.rs 作为 Future 轮询的运行时(Hysteria 2 的运行时和 TUN 的 UDP 运行时;drive 使用 stream 模式)在完成这么多单元后让出
INLINE_EFFECTS 4 concepts/src/core.rs effect 列表溢出到堆之前内联存放的 effect 数
MAX_SUBS 64 app/src/outbound/udp_fanout.rs 每个 UDP 关联打开的子链路数
HTTP_BUF、SOCKS_BUF、TROJAN_BUF、VLESS_BUF 16 KiB app/src/outbound/mod.rs 代理客户端运行时缓冲区
VMESS_BUF、SS2022_BUF;SS_BUF 32 KiB;20 KiB app/src/outbound/mod.rs 代理客户端运行时缓冲区

每条连接的内存由缓冲区大小决定。一个服务端运行时恰好分配三个 BUF_SIZE 数组(读、暂存、scratch),因此一条 VMess 连接的缓冲区占 96 KiB,一条 HTTP 连接占 192 KiB。代理出站再加上其客户端运行时的两个缓冲区(暂存和读),例如 VMess 为 2 × 32 KiB。每打开一个出站,还会多一个 KeyWaker Arc 和一个装箱的连接 future。

  • concepts/tests/runtime.rs 用简化的协议核心驱动运行时:打开、中继、半关闭、连接失败、截止时间、held 字节和背压。stream_mode_reports_each_unit_of_work 固定了 drive 所依赖的 showing_progress 模式,deadline_event_lets_the_core_time_out 和 deadline_is_armed_against_the_tokio_clock 固定了 Effect::SetDeadline 背后的计时器。
  • concepts/tests/client.rs 覆盖 ProxyClientRuntime 以及 ProxyClientConnector 的 Connected 语义。
  • protocols/tests/unit/*/core.rs 手动驱动每个协议核心,大多数通过 CoreHarness:它何时打开,在 Connected 和 ConnectFailed 时暂存什么。
  • protocols/tests/pipeline/ 在真实运行时中运行每个服务端协议核心,通过回环 socket 与其客户端 codec 对接(support::pipeline::serve_runtime、client);transports.rs 覆盖每种入站传输层,包括一条 gRPC 连接上的多条 stream。
  • app/tests/integration/ 启动真实的 etemenanki-app 二进制:e2e_route_context.rs、e2e_sniff.rs 和 e2e_udp_route.rs 端到端地测试第 6、7 步以及 UDP 扇出。
  • app/tests/integration/e2e_unix.rs(socks_over_a_unix_socket_relays_and_cleans_up、http_connect_over_a_unix_socket_relays)覆盖跳过传输层的 Unix 路径;e2e_hysteria_inbound.rs 和 e2e_tun.rs 覆盖拥有自己运行循环的两个入站。
  • app/tests/unit/serve.rs → accept_error_backoff_classification 固定了 accept 错误的分类;它是 app/src/serve.rs 唯一的单元测试,因此 app 的 drive 看门狗没有专门测试。katana 的 tests/unit/serve.rs → a_silent_client_is_dropped_at_the_handshake_deadline 固定了 katana 自己的看门狗,其形态不同(整个握手只有一个超时,而不是每步一个)。

各层测试的运行方式见 测试。