跳转到内容

出站、UDP fan-out 与负载均衡器

源码文件:72 个 · 核对版本 Etemenanki 555b7df
  • Etemenanki/supervisor/src/build/outbound.rs
  • Etemenanki/supervisor/src/topology/outbound/mod.rs
  • Etemenanki/supervisor/src/topology/outbound/proxy.rs
  • Etemenanki/supervisor/src/topology/outbound/freedom.rs
  • Etemenanki/supervisor/src/topology/outbound/udp_fanout.rs
  • Etemenanki/supervisor/src/topology/balancer.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/topology/plane.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/topology/spec_plan/outbound.rs
  • Etemenanki/supervisor/src/topology/spec_plan/transport.rs
  • Etemenanki/supervisor/src/topology/spec_plan/route.rs
  • Etemenanki/supervisor/src/topology/spec_plan/plan.rs
  • Etemenanki/supervisor/src/build/validate.rs
  • Etemenanki/supervisor/src/build/dns.rs
  • Etemenanki/supervisor/src/build/apply.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/track/mod.rs
  • Etemenanki/supervisor/src/track/metered.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/socks/codec.rs
  • Etemenanki/protocols/src/http/codec.rs
  • Etemenanki/protocols/src/trojan/codec.rs
  • Etemenanki/protocols/src/trojan/protocol.rs
  • Etemenanki/protocols/src/vless/codec.rs
  • Etemenanki/protocols/src/vless/protocol.rs
  • Etemenanki/protocols/src/vmess/codec.rs
  • Etemenanki/protocols/src/ss_legacy/codec.rs
  • Etemenanki/protocols/src/ss_2022/codec.rs
  • Etemenanki/protocols/src/ss_2022/crypto.rs
  • Etemenanki/protocols/src/ss_aead.rs
  • Etemenanki/protocols/src/helpers/address.rs
  • Etemenanki/protocols/src/helpers/address_family.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/transports/keepalive.rs
  • Etemenanki/protocols/src/transports/tls/config.rs
  • Etemenanki/protocols/src/hysteria/config.rs
  • Etemenanki/protocols/src/hysteria/connector.rs
  • Etemenanki/protocols/src/hysteria/connection.rs
  • Etemenanki/protocols/src/hysteria/slot.rs
  • Etemenanki/protocols/src/wireguard/config.rs
  • Etemenanki/protocols/src/wireguard/connector.rs
  • Etemenanki/protocols/src/wireguard/slot.rs
  • Etemenanki/protocols/src/wireguard/device.rs
  • Etemenanki/environment/src/dial/mod.rs
  • Etemenanki/environment/src/dial/tcp.rs
  • Etemenanki/environment/src/dial/udp.rs
  • Etemenanki/environment/src/dial/socket.rs
  • Etemenanki/app/src/lower.rs
  • Etemenanki/app/src/subscribe.rs
  • Etemenanki/ffi/src/platform.rs
  • Etemenanki/supervisor/tests/unit/balancer.rs
  • Etemenanki/supervisor/tests/unit/plane.rs
  • Etemenanki/supervisor/tests/unit/plan.rs
  • Etemenanki/supervisor/tests/unit/validate.rs
  • Etemenanki/supervisor/tests/socket_policy.rs
  • Etemenanki/supervisor/tests/hot_swap.rs
  • Etemenanki/supervisor/tests/tracking.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/app/tests/integration/e2e_udp_route.rs
  • Etemenanki/app/tests/integration/e2e_balancer.rs
  • Etemenanki/app/tests/integration/e2e_route_context.rs
  • Etemenanki/app/tests/integration/e2e_xray.rs
  • Etemenanki/app/tests/integration/e2e_xray_vmess.rs
  • Etemenanki/app/tests/integration/e2e_hysteria.rs
  • Etemenanki/app/tests/integration/e2e_wg.rs

出站是被路由的流离开本进程的地方。supervisor(监管器)在一次应用(apply)需要时为每个 OutboundSpec 构建一个 handler(处理器),只要它的 spec(期望状态)和 DNS spec 保持不变就一直保留它,并通过两个方法在它上面打开流:一个用于 TCP 流,一个用于 UDP 流。本页介绍 etemenanki-supervisor 中的这一层:build_outbound 以及交给每个协议的解析器和 socket 策略,封闭的 Outbound 枚举及其返回的 stream 与数据报类型,代理客户端及其缓冲区,SOCKS、freedom、blackhole、Hysteria 2 和 WireGuard 出站,逐包的 UDP fan-out(扇出),负载均衡器,以及出站打开的每个 socket 上的 socket 策略。

本页面向这样的贡献者:添加出站协议、修改 UDP 关联的路由方式,或者改动负载均衡。流如何经由 plane(数据平面)到达目标(Plane、Target、AppConnector::connect),以及排空对已打开的流做什么,见 plane:为每个流选路。每个代理出站所运行的客户端运行时见客户端 codec 与客户端运行时。运维者视角下的同一组功能见用户指南的出站和负载均衡器。

组件 文件 → 符号 负责
构建 supervisor/src/build/outbound.rs → build_outbound、transport 把一个 OutboundSpec 变成一个 Outbound:每个流的 codec 闭包、传输层 connector、解析器和 socket 策略。
分发 supervisor/src/topology/outbound/mod.rs → Outbound 把 TCP 流作为字节流打开(connect_stream),把 UDP 流作为数据报 link 打开(connect_datagram),并在协议不支持 UDP 时拒绝 UDP。
代理客户端 supervisor/src/topology/outbound/proxy.rs → ProxyClient 每个出站一个位于互斥锁之后的 ProxyClientConnector,以及封闭的 OutboundStream、OutboundDatagram 和 BlackholeLink 类型。
SOCKS supervisor/src/topology/outbound/mod.rs → SocksOutbound 通过 ProxyClient 执行 CONNECT,以及把 UDP ASSOCIATE 作为一个独立的 link,建立在一条控制 stream 之上。
直连 supervisor/src/topology/outbound/freedom.rs → FreedomConnector、ResolvingUdp 解析并连接 TCP;按地址族绑定 UDP socket,并解析域名目的地。
UDP fan-out supervisor/src/topology/outbound/udp_fanout.rs → FanOutLink 在当前 plane 上路由 UDP 关联的每个包,每个出站版本保留一个子 link,并合并各方的回复。
负载均衡器 supervisor/src/topology/balancer.rs → Balancer、Member、Strategy 每个成员的健康状态,每个成员一个探测任务,以及按流或按包选择成员。

它交给别处的事:

  • 路由。 Plane::route、Target、Guarded 和 AppConnector::connect 见 plane:为每个流选路。本页只在 fan-out 和负载均衡器需要的范围内使用 Target::resolve、Target::connect_datagram 和 Guarded。
  • 版本、复用与排空。 一次应用构建、复用或排空哪些出站由计划决定;见规划并应用变更。
  • 计量、强制关闭(kill)与限速。 每个 stream 和每个 fan-out 子 link 都是一个被跟踪的流;见流跟踪、统计与限速。
  • DNS 服务与解析器。 Outbound::Dns、DnsUdpLink、DnsTcpStream、RoutedDialer 和 Dns 解析器对见名称解析与 DNS 服务。
  • spec 规则。 validate 拒绝的一切见校验与应用错误;本页只列出负载均衡器的检查。
  • 报文格式。 每个 codec 和客户端见各自的协议页,例如 SOCKS、Hysteria 2:协议与客户端 和 WireGuard。
  1. plan(supervisor/src/topology/spec_plan/plan.rs)为每个出站 tag 分配一个版本;两者合起来是 OutboundId { tag, version }(supervisor/src/entity/id.rs)。spec 未变的出站是运行中版本上的 Step::Reuse。其他情况都是下一个版本上的 Step::Build,旧版本得到一个 Step::Drain。DNS spec 变化时,除 blackhole 外的每个出站都会重建,因为其他出站都持有由它构建的解析器(uses_dns)。已移除 tag 的版本保留在 RunningState::versions 中,所以重新加回的 tag 不会复用一个可能仍有流在排空的版本。
  2. Actor::prepare(supervisor/src/supervisor.rs)为每个要构建的出站调用 build_outbound(outbound, &dns, &self.socket),并用 Target::outbound(id, built) 把结果包进一个 Arc。被复用的出站就是运行中的那个 Arc<Target> 本身,所以它保留自己持有的东西:它的 ProxyClient、它的 QUIC 连接或它的隧道。
  3. 接下来构建负载均衡器,因为成员是同一次应用中构建或复用的出站的 Arc<Target>(见多次应用之间的状态)。
  4. commit 发布新的 plane,然后停止被重建或移除的负载均衡器的探测,并为新的负载均衡器启动探测。

构建失败是一个 ApplyError::Build,打印为 building outbound <tag>@v<n> failed: <error>,整个应用被拒绝,运行中的状态保持原样。etemenanki-app 把它打印在 configuration invalid: 之后。对于 ca_file 中没有证书的 TLS 出站,本文核对版本的二进制打印:

configuration invalid: building outbound up@v1 failed: no certificate in CA PEM bundle

check(即 app 的 --test)运行 prepare 但不绑定。它构建每个出站和每个负载均衡器,所以构建错误都能被发现;但出站在构建时不打开任何 socket,也不启动任何探测。出站打开的第一个 socket 属于它的第一个流。

supervisor/src/build/outbound.rs
pub(crate) fn build_outbound(
spec: &OutboundSpec,
dns: &Dns,
socket: &SocketOptions,
) -> io::Result<Outbound>;
pub(crate) fn obfs(spec: &ObfsSpec) -> Obfs;
fn tcp(upstream: &ProxyUpstream) -> Destination;
fn udp(server: &Destination) -> Destination;
fn transport(
spec: &OutboundTransportSpec,
dns: &Dns,
socket: &SocketOptions,
) -> io::Result<TransportConnector>;

Dns(supervisor/src/build/dns.rs)带有两个解析器。servers 查询出站要拨号的代理服务器。destinations 查询本进程自己要连接的目的地:经由 freedom 访问的,以及 WireGuard 隧道内访问的。DnsSpec::Single 时两者是同一个解析器。DnsSpec::Split 时,servers 直接询问 pre_proxy 服务器,destinations 经由路由表询问 through_proxy 服务器;见名称解析与 DNS 服务。

socket 是 supervisor 的 socket 策略,会交给构建协议时用到的每个拨号器(见每个出站 socket 上的 socket 策略)。

每种 spec 变成什么:

OutboundProtocolSpec Outbound 变体 构建方式
Freedom { address_family } Freedom(FreedomConnector) FreedomConnector::new(Dialer::new(socket), dns.destinations, address_family)
Blackhole Blackhole 无
Socks { upstream, account } Socks(Box<SocksOutbound>) SocksOutbound::new(transport(..), UdpDialer::new(socket), tcp(upstream), auth)
Http { upstream, account } Http(ProxyClient<HTTP_BUF, HttpConnect, NoUdp>) make 构建 HttpConnect::new(&flow.destination, auth)
Trojan { upstream, password } Trojan(ProxyClient<TROJAN_BUF, TrojanStream, TrojanDatagram>) password_hash(password) 只算一次;make 构建 TrojanStream::new(&hash, &dest) 或 TrojanDatagram::new(&hash)
Vless { upstream, id } Vless(ProxyClient<VLESS_BUF, VlessStream, VlessDatagram>) make 构建 VlessStream::new(&uuid, &dest) 或 VlessDatagram::new(&uuid, &dest)
Vmess { upstream, id, security } Vmess(ProxyClient<VMESS_BUF, VMessStream, VMessDatagram>) make 构建 VMessStream::new(uuid, security, true, &dest) 或 VMessDatagram::new(..);true 即 global_padding
Shadowsocks { upstream, method, password } Shadowsocks(ProxyClient<SS_BUF, SsStream, NoUdp>) evp_bytes_to_key(password, method.key_len()) 只算一次;make 构建 SsStream::new(method, key.clone(), &dest)
Ss2022 { upstream, method, identity_psks, psk } Ss2022(ProxyClient<SS2022_BUF, Ss2022Stream, NoUdp>) make 构建 Ss2022Stream::new(method, psk.clone(), identity.clone(), &dest)
Hysteria2(spec) Hysteria2(Hy2Connector) Hy2Connector::with_address_family(Hy2Config { .. }, address_family).with_resolver(dns.servers)
Wireguard(spec) Wireguard(WgConnector) WgConnector::with_address_family(WgConfig { .. }, address_family).with_resolver(dns.destinations)

每种出站使用哪个解析器、打开哪些 socket:

协议 代理服务器或端点由谁解析 目的地由谁解析 在 socket 策略下打开的 socket
freedom 无 destinations,按出站的 address_family TCP 连接;每个地址族一个 UDP socket
blackhole 无 无 无
socks servers,按传输层的 address_family 上游 传输层之下的 TCP 连接;UDP 中继 socket
http、trojan、vless、vmess、shadowsocks(传统与 2022 两类) servers,按传输层的 address_family 上游 传输层之下的 TCP 连接
hysteria2 servers,按出站的 address_family 服务器 QUIC 的 UDP socket
wireguard servers,按 endpoint_address_family destinations,按 address_family 以及隧道自身地址的地址族 隧道的 UDP socket

表中的三个细节:

  • 密钥只派生一次。 Trojan 密码的哈希和 Shadowsocks 主密钥在构建出站时计算,make 闭包捕获其结果。每个流只构建自己的 codec。
  • tcp 与 udp 固定网络类型。 tcp(upstream) 复制代理服务器的 Destination 并设 network = DialNetwork::Tcp;udp(server) 对 Hysteria 2 服务器和 WireGuard 端点做同样的事,设为 DialNetwork::Udp。
  • obfs 把 spec 的 ObfsSpec::Salamander { psk } 转换为 protocols crate 的 Obfs::Salamander { psk }。

transport(spec, dns, socket) 把上游的 StreamShape 映射为一个 TransportKind。tls(alpn) 是一个基于 OutboundTransportSpec::tls 的闭包,返回 Option<ClientConfig>:spec 带有 TLS 材料时为 Some,否则为 None。

StreamShape 构建的 TransportKind 是否叠加 TLS 缺字段时的错误
Tcp TransportKind::Tcp 从不 无
Tls TransportKind::Tls(config),来自 tls(Alpn::None) 总是,无 ALPN missing tls
Ws { path, host, .. } TransportKind::ws(host, path, tls(Alpn::Http1)?) 设置了 OutboundTransportSpec::tls 时叠加,ALPN 为 http/1.1 missing ws host
Grpc { service, authority, .. } TransportKind::grpc(authority, service, tls(Alpn::Http2)?) 设置了 OutboundTransportSpec::tls 时叠加,ALPN 为 h2 missing grpc authority

transport 忽略 shape 自身的 tls 标志:当且仅当设置了 OutboundTransportSpec::tls 时才叠加 TLS,而校验保证两者一致(outbound_transport 拒绝 shape.uses_tls() != tls.is_some() 的 spec)。三个 missing … 错误的 kind 都是 InvalidInput,对通过校验的 spec 不会出现:出站的传输层 spec 没有空缺,因为前端程序会填好 WebSocket host、gRPC authority 和 SNI,校验也会检查它们都在。

TransportKind::grpc 总是构建 GrpcMode::Gun(每条 gRPC 消息一个 Hunk),打开时使用 user-agent DEFAULT_USER_AGENT,这是一个桌面版 Chrome 的字符串(protocols/src/transports/grpc/settings.rs)。protocols crate 还提供了用于 MultiHunk 分帧的 multi(),以及用于修改或去掉该头部的 user_agent(),但 supervisor 两者都不调用,所以每个 gRPC 出站都使用 gun 分帧和默认 user agent。

TLS 配置是 ClientConfig::with_verify_mode(&tls.server_name, tls.verify.clone(), alpn)(protocols/src/transports/tls/config.rs)。它构建一个最低版本为 TLS 1.2 的 OpenSSL SslConnector,并按 VerifyMode 设置校验:

VerifyMode 信任 是否检查主机名
System 默认校验路径(系统根证书) 是
CustomCa(pem) 默认校验路径,加上证书包中的每个证书 是
Insecure 无:SslVerifyMode::NONE 否

CustomCa 证书包在这里、即构建出站时解析。无法解析的证书包以 OpenSSL 自己的错误文本失败(kind 为 Other),解析后没有证书的则以 no certificate in CA PEM bundle 失败(InvalidInput)。

结果是 TransportConnector::new(kind, Dialer::new(socket.clone()), dns.servers.clone(), spec.address_family)。TransportConnector::dial(protocols/src/transports/connect.rs)按以下步骤执行:

  1. 它以 Unsupported a proxy transport carries no datagrams of its own 拒绝 UDP 目的地。代理客户端交给它的总是 tcp(upstream),代理的 UDP 在它的 stream 内部承载。
  2. 它用 servers 按传输层的 address_family 解析代理服务器(destination_to_socketaddrs)。
  3. 它连接第一个应答的地址(TcpDialer::connect_any)。
  4. 它用 set_keepalive(protocols/src/transports/keepalive.rs)为连接设置 TCP keepalive:静默 TCP_KEEPALIVE_IDLE(120 秒)后发出第一个探测包,之后每 TCP_KEEPALIVE_INTERVAL(30 秒)一个,连续 TCP_KEEPALIVE_RETRIES(3)个探测包无应答就放弃连接。常量上的注释给出了原因:Linux 要等两小时才发第一个探测包,而代理的对端经常不发 FIN 就消失,例如移动网络的 NAT 重新绑定之后。平台拒绝这些选项不算错误;这种失败以 debug 级别记录为 could not enable TCP keepalive: <error>。
  5. 它把 stream 包进传输层。

传输层见传输层:TCP 与 TLS和传输层:WebSocket 与 gRPC。

supervisor/src/topology/outbound/mod.rs
const HTTP_BUF: usize = 16 * 1024;
const SOCKS_BUF: usize = 16 * 1024;
const TROJAN_BUF: usize = 16 * 1024;
const VLESS_BUF: usize = 16 * 1024;
const VMESS_BUF: usize = 32 * 1024;
const SS_BUF: usize = SsStream::BUF_SIZE;
const SS2022_BUF: usize = Ss2022Stream::BUF_SIZE;
pub enum Outbound {
Freedom(FreedomConnector),
Blackhole,
Socks(Box<SocksOutbound>),
Http(ProxyClient<HTTP_BUF, HttpConnect, NoUdp>),
Trojan(ProxyClient<TROJAN_BUF, TrojanStream, TrojanDatagram>),
Vless(ProxyClient<VLESS_BUF, VlessStream, VlessDatagram>),
Vmess(ProxyClient<VMESS_BUF, VMessStream, VMessDatagram>),
Shadowsocks(ProxyClient<SS_BUF, SsStream, NoUdp>),
Ss2022(ProxyClient<SS2022_BUF, Ss2022Stream, NoUdp>),
Wireguard(WgConnector),
Hysteria2(Hy2Connector),
Dns(etemenanki_protocols::dns::serve::Service),
}
pub type StreamFuture = Pin<Box<dyn Future<Output = io::Result<OutboundStream>> + Send>>;
pub type DatagramFuture = Pin<Box<dyn Future<Output = io::Result<OutboundDatagram>> + Send>>;
impl Outbound {
pub fn connect_stream(&self, flow: Flow) -> StreamFuture;
pub fn connect_datagram(&self, flow: Flow) -> DatagramFuture;
}

不要把它和 etemenanki_concepts::link::Outbound<S, D> 混淆,后者是每个 Connector 返回的“stream 或数据报”二选一结果;下文代码把它称为 link::Outbound。Flow 即 etemenanki_protocols::flow::Flow<Principal>(supervisor/src/topology/flow.rs)。

这个枚举是有意封闭的。一个 Target 拥有一个 Outbound,而服务端运行时按值保存它打开的每个出站,所以每一类流都需要一个具体类型。两个方法都接受 &self,因为路由到同一目标的所有流会同时在同一个值上拨号。每个分支同步构建自己的拨号 future 并装箱;在 future 第一次被 poll 之前不会触碰任何 socket。

变体 connect_stream connect_datagram
Freedom 克隆 connector 并拨号;stream 成为 OutboundStream::Tcp 克隆 connector 并拨号;ResolvingUdp link 被装箱进 OutboundDatagram
Blackhole 立即就绪:OutboundStream::Blackhole 立即就绪:基于 BlackholeLink 的 OutboundDatagram
Socks SocksOutbound.tcp.connect(flow),然后 proxy_stream SocksOutbound::connect_datagram(),即一次 UDP ASSOCIATE;不使用该流
Http ProxyClient::connect,然后 proxy_stream 失败:http carries no datagrams
Trojan、Vless、Vmess ProxyClient::connect,然后 proxy_stream ProxyClient::connect,然后 proxy_datagram
Shadowsocks、Ss2022 ProxyClient::connect,然后 proxy_stream 失败:shadowsocks carries no datagrams、shadowsocks-2022 carries no datagrams
Wireguard 克隆 connector 并拨号;stream 成为 OutboundStream::Wg 克隆 connector 并拨号;link 被装箱进 OutboundDatagram
Hysteria2 克隆 connector 并拨号;stream 成为 OutboundStream::Hy2 克隆 connector 并拨号;link 被装箱进 OutboundDatagram
Dns DnsTcpStream::new(service) 装箱进 OutboundStream::Proxy DnsUdpLink::new(service) 装箱进 OutboundDatagram

这些“no datagrams”错误来自 no_udp(what),kind 为 Unsupported。FreedomConnector、WgConnector 和 Hy2Connector 每次拨号都会克隆,因为它们的 connect 接受 &mut self;克隆的开销很小。WgConnector 的克隆共享它的隧道槽(Arc<tokio::sync::Mutex<DeviceSlot>>),Hy2Connector 的克隆共享它的连接槽(Arc<parking_lot::Mutex<ConnSlot>>),所以克隆永远不会打开第二条隧道或连接。

有七个协议是基于 TransportConnector 的客户端:HTTP、SOCKS CONNECT、Trojan、VLESS、VMess、Shadowsocks 和 Shadowsocks 2022。它们都经过同一个包装:

supervisor/src/topology/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>;
}

为什么要互斥锁。 ProxyClientConnector::connect 接受 &mut self:Make 是 FnMut,内层的 TransportConnector 也是一个 connect 接受 &mut self 的 Connector。出站是共享的,所以 ProxyClient 把 connector 放在一个 parking_lot::Mutex 里。connect 加锁,调用内层的 connect,在它返回时解锁。在锁内,connector 运行 make,构建 ProxyClientRuntime(codec 加上两个 BUF 字节的装箱缓冲区),并向 TransportConnector 要它的拨号 future,后者把 connector 克隆进一个装箱的 async 块。锁内没有任何 I/O,也没有 .await,所以经过同一出站的流只在构建 future 时串行,从不在网络上串行。

codec 按流选择。 每个 make 闭包返回 link::Outbound::Stream(codec) 或 link::Outbound::Datagram(codec)。Trojan、VLESS 和 VMess 依据 flow.destination.network == DialNetwork::Udp 选择。HTTP、SOCKS、Shadowsocks 和 Shadowsocks 2022 总是返回 stream,NoUdp 充当它们的数据报 codec。NoCodec(concepts/src/core.rs)是一个没有值的(uninhabited)枚举 Never(Infallible, PhantomData<..>),所以它不可能有任何值,它的 codec 方法也永远不会被调用。NoUdp 协议永远不会带着 UDP 流到达 make,因为 connect_datagram 先拒绝了它;SOCKS UDP 根本不使用 ProxyClient。

拨号 future。 只有在上游已拨通、codec 的握手已完成、握手字节已 flush 之后,ProxyClientConnecting 才会完成。因此上游拒绝表现为 future 返回的 Err,服务端核心看到的是 ConnectFailed,而不是一个稍后才失败的 stream。随后两个适配函数把结果折叠成封闭类型:

supervisor/src/topology/outbound/proxy.rs
pub fn proxy_stream<const BUF: usize, S, D>(
opened: link::Outbound<
ProxyClientRuntime<BUF, S, TransportConnector>,
ProxyClientRuntime<BUF, D, TransportConnector>,
>,
) -> io::Result<OutboundStream>;
pub fn proxy_datagram<const BUF: usize, S, D>(
opened: link::Outbound<
ProxyClientRuntime<BUF, S, TransportConnector>,
ProxyClientRuntime<BUF, D, TransportConnector>,
>,
) -> io::Result<OutboundDatagram>;

类型不对的运行时会得到 a TCP flow was dialed as datagrams 或 a UDP flow was dialed as a stream(kind 为 Other)。实际中两者都不会发生,因为 make 依据的 network 字段,正是在 connect_stream 与 connect_datagram 之间做选择的同一个字段。

每个 BUF 常量决定 ProxyClientRuntime 两个缓冲区的大小:发往上游的暂存缓冲区,以及从上游读取的读缓冲区。常量上的注释给出了规则:每个协议要解开的最大线上帧,加上它的 codec 预留的空间。

出站 常量 BUF codec codec 的 STAGING_RESERVE 每个流的缓冲区
http HTTP_BUF 16,384 HttpConnect REQUEST_MAX = 1,024 32 KiB
socks(CONNECT) SOCKS_BUF 16,384 SocksConnect 528 32 KiB
trojan TROJAN_BUF 16,384 TrojanStream、TrojanDatagram REQUEST_HEADER_MAX = 56 + 2 + 1 + 259 + 2 = 320 32 KiB
vless VLESS_BUF 16,384 VlessStream、VlessDatagram REQUEST_HEADER_MAX = 1 + 16 + 1 + 1 + 259 = 278 32 KiB
vmess VMESS_BUF 32,768 VMessStream、VMessDatagram HEADER_MAX(358)向上取整到 64 的倍数 = 384 64 KiB
shadowsocks(传统 AEAD) SS_BUF = SsStream::BUF_SIZE 20,480 SsStream 32 + 34 + 259 + 61 = 386 40 KiB
shadowsocks(2022- 方法) SS2022_BUF = Ss2022Stream::BUF_SIZE 65,569 Ss2022Stream 2,048 131,138 字节

259 是 AddressCodec::MAX_LEN(类型、长度、255 字节的名称、端口),34 是 AEAD 的 RECORD_OVERHEAD(一个加密的两字节长度及其 16 字节 tag,加上负载的 16 字节 tag)。

两个 Shadowsocks 的大小来自它们的 codec:

  • SsStream::BUF_SIZE 容纳 salt 和一个 MAX_PAYLOAD(0x3FFF 字节)的 chunk 及其开销,还有余量。
  • Ss2022Stream::BUF_SIZE 是 MAX_RECORD_LEN = MAX_PACKET_SIZE(0xFFFF)+ RECORD_OVERHEAD = 65,569。ss-rust 和 sing 会把一次最多 0xFFFF 字节的读取整个封装进一条记录,其他帧都更小,所以这么大的缓冲区总能完整容纳下一帧。

ProxyClientRuntime::new 断言 BUF_SIZE > STAGING_RESERVE(BUF_SIZE must exceed the codec's STAGING_RESERVE or nothing can be sealed);上表每一组都满足。缓冲区限制了每个方向上单帧的大小:

  • 来自上游。 读缓冲区装不下的帧会使流以 InvalidData upstream frame larger than the client runtime's buffer 失败。
  • 发往上游。 stream 写入分段封装,所以不会溢出。数据报则整体封装:负载加上 codec 的 STAGING_RESERVE 超过 BUF 的 UDP 包(例如 Trojan 中负载超过 16,064 字节)会使这次发送以 InvalidInput frame larger than the client runtime's buffer 失败。在 fan-out 中,发送失败会丢弃该包和这个子 link(见发送)。

缓冲区机制和运行时的其他错误见客户端 codec 与客户端运行时;一个流可能遇到的错误列在错误中。

Section titled “OutboundStream、OutboundDatagram 与 BlackholeLink”
supervisor/src/topology/outbound/proxy.rs
pub trait Stream: AsyncRead + AsyncWrite + Send {}
impl<T: AsyncRead + AsyncWrite + Send> Stream for T {}
pub type ProxyStream = Pin<Box<dyn Stream>>;
pub enum OutboundStream {
Tcp(TcpStream),
Proxy(ProxyStream),
Wg(WgStream),
Hy2(Hy2Stream),
Blackhole,
}
pub struct OutboundDatagram(Box<dyn DatagramLink<Addr = Destination> + Send>);
impl OutboundDatagram {
pub fn new<D: DatagramLink<Addr = Destination> + Send + 'static>(link: D) -> Self;
}
pub struct BlackholeLink;
  • OutboundStream 通过委托给它的变体来实现 AsyncRead 和 AsyncWrite。代理客户端的类型随协议和缓冲区大小而不同,所以它们被装箱进 Proxy;DNS 服务的 DnsTcpStream 也使用这个变体。普通 TCP、WireGuard 和 Hysteria 2 的 stream 有各自的变体,免去装箱。
  • OutboundStream::Blackhole 读取时立即得到 EOF,写入时整块接受;flush 和 shutdown 立即成功。
  • OutboundDatagram 把每个数据报 link 装箱,因为 fan-out 要把不同出站的 link 并排保存。
  • BlackholeLink 把每次 poll_send_to 报告为全部发送,每次 poll_recv_from 永远返回 Pending。

stream 或 link 在到达服务端运行时之前会被包装两层:先是 Guarded(它的出站版本被排空策略关闭后,它就会失败,见 plane:为每个流选路),然后是 Metered 或 MeteredDatagram(见流跟踪、统计与限速)。TCP 流是 Metered<Guarded<OutboundStream>>;fan-out 子 link 是 MeteredDatagram<Guarded<OutboundDatagram>>。

supervisor/src/topology/outbound/mod.rs
pub struct SocksOutbound {
transport: TransportConnector,
/// Binds the socket a `UDP ASSOCIATE` relays through.
udp: UdpDialer,
server: Destination,
auth: Option<(CompactString, CompactString)>,
tcp: ProxyClient<SOCKS_BUF, SocksConnect, NoUdp>,
}
impl SocksOutbound {
pub fn new(
transport: TransportConnector,
udp: UdpDialer,
server: Destination,
auth: Option<(CompactString, CompactString)>,
) -> Self;
fn connect_datagram(&self) -> DatagramFuture;
}

TCP 流经过 tcp,这是一个普通的 ProxyClient,它的 make 构建 SocksConnect::new(&flow.destination, auth)。UDP 不适合客户端运行时,因为 SOCKS5 关联由两条连接组成:一条必须保持打开的控制 stream,以及一个发往服务器中继的 UDP socket。connect_datagram 自己构建这个 link:

  1. 它用 transport.dial(&server) 拨一条控制 stream,所以出站的传输层(TLS、WebSocket、gRPC)及其 TCP keepalive 都作用于控制 stream。
  2. 它运行 SocksUdpLink::associate(control, auth, bind):先是方法协商,设置了 auth 时再进行一轮用户名密码认证,然后是 UDP ASSOCIATE。方法请求只提供一种方法(encode_method_request):设置了 auth 时为用户名密码(0x02),否则为无认证(0x00)。每一轮都以 512 字节为块读取控制 stream,读进同一个缓冲区,直到它的解析器(parse_method_reply、parse_userpass_reply、parse_reply)拿到完整的回复;解析器错误(例如版本字节不对)原样返回。
  3. bind 闭包是 udp.bind(AddressFamily::of(relay)):在 socket 策略下、按中继的地址族、在临时端口上绑定一个 socket。UdpDialer 把 IPv6 socket 设为 v6-only。
protocols/src/socks/udp_link.rs
impl<S> SocksUdpLink<S>
where
S: AsyncRead + AsyncWrite + Unpin,
{
pub async fn associate(
mut control: S,
auth: Option<(&str, &str)>,
bind: impl FnOnce(&SocketAddr) -> io::Result<UdpSocket>,
) -> io::Result<Self>;
pub fn relay(&self) -> SocketAddr;
}
sequenceDiagram
  participant F as FanOutLink
  participant O as SocksOutbound
  participant T as TransportConnector
  participant S as SOCKS5 服务器
  F->>O: connect_datagram
  O->>T: 拨号到服务器
  T->>S: TCP 连接,然后传输层握手
  O->>S: 方法请求
  S-->>O: 选定的方法
  opt 设置了 auth
    O->>S: 用户名和密码
    S-->>O: 状态
  end
  O->>S: UDP ASSOCIATE,指定 0.0.0.0 端口 0
  S-->>O: 指明中继地址的回复
  O->>O: 按中继的地址族绑定一个 UDP socket
  O-->>F: 基于 SocksUdpLink 的 OutboundDatagram
  F->>O: poll_send_to
  O->>S: 中继头和负载,发往中继
  S-->>O: 来自中继的数据报
  O-->>F: 负载,以及头部指明的对端

请求不指明来源:encode_request(CMD_UDP_ASSOCIATE, None) 写入 0.0.0.0:0,这正是 RFC 1928 规定客户端在不知道自己来源时的做法。socket 要等回复指明中继后才绑定,而在 NAT 之后本地地址本来就不对。检查来源的服务器会把关联限定在控制连接的来源地址和第一个数据报的端口上,所以数据报必须从服务器在控制连接上看到的那个 IP 发出。

关联建立后,link 的行为如下:

  • 每次 poll_send_to 都按该包自己的目的地加上 SOCKS5 UDP 头(encode_udp_packet_into,写入一个复用的、初始容量为 2,048 的 scratch 缓冲区),所以一个关联可以到达多个对端。
  • poll_recv_from 读入一个装箱的 RECV_BUF(64 KiB)缓冲区。它跳过来源不是中继的数据报(endpoint(from) != endpoint(self.relay),比较的是规范化的 IP 和端口,所以在双栈 socket 上以 IPv4 映射地址收到的 IPv4 中继仍然算数),以及头部无法解析的数据报。它把负载复制到调用方的缓冲区,装不下的负载会被截断,与 UDP 的行为一致,并返回头部指明的对端。
  • 两个方向都会先把控制 stream 读空到一个 256 字节的 sink 中;服务器在上面发送的任何东西都被丢弃。控制 stream 读到 EOF 就结束关联:这次调用和之后的每次调用都以 BrokenPipe socks: the control connection closed 失败。控制 stream 上的读错误原样返回一次,并把 stream 标记为已关闭,所以之后的每次调用都以同样的 BrokenPipe 错误失败。
  • 握手以 PermissionDenied(auth method not supported、server rejects account)、ConnectionRefused(server rejects request: <status>)、Unsupported(socks: the relay address is a domain)或 UnexpectedEof(socks: server closed during the handshake)失败。

报文格式和关联的服务端见 SOCKS。

supervisor/src/topology/outbound/freedom.rs
#[derive(Clone)]
pub struct FreedomConnector {
dialer: Dialer,
resolver: Resolver,
strategy: AddressFamilyStrategy,
}
impl FreedomConnector {
pub fn new(dialer: Dialer, resolver: Resolver, strategy: AddressFamilyStrategy) -> Self;
}
impl Default for FreedomConnector; // Dialer::default(), Resolver::default(), Auto
impl Connector<Flow> for FreedomConnector {
type Stream = TcpStream;
type Datagram = ResolvingUdp;
type Future = DialFuture;
fn connect(&mut self, flow: Flow) -> DialFuture;
}
pub struct ResolvingUdp {
socket: DualStackUdp,
resolver: Resolver,
strategy: AddressFamilyStrategy,
// 私有的解析状态
}
impl DatagramLink for ResolvingUdp {
type Addr = Destination;
// poll_send_to, poll_recv_from
}

connect 把 connector 克隆进一个装箱的 dial(flow.destination)。解析器是本次应用的 destinations 解析器,策略是出站的 address_family。Default 基于系统解析器和默认 socket 选项构建一个实例;supervisor 中没有任何地方使用它。

TCP。 destination_to_socketaddrs(&dest, strategy, &resolver) 解析目的地并对候选地址排序:auto 保持解析器给出的顺序,ipv4_only 和 ipv6_only 做过滤,prefer_ipv4 和 prefer_ipv6 做稳定排序,让另一个地址族留作后备。然后 dialer.tcp.connect_any(&addrs) 按顺序逐个尝试这些地址,每次尝试受 DEFAULT_CONNECT_TIMEOUT(10 秒)限制。如果每次尝试都失败,错误是 ConnectionRefused failed to connect to any address (<addr>: <error>; …),列出每个地址及其各自的错误。尝试是有意按顺序进行的(不是 Happy Eyeballs):双栈主机经常把一个名称解析到自己无法到达的地址,而保留每次失败可以避免丢失原因。freedom 不调用 set_keepalive,所以直连 TCP 连接只有 socket 策略设置的 keepalive(见每个出站 socket 上的 socket 策略)。

UDP。 dialer.udp.bind_dual(families) 每个地址族最多绑定一个 socket。私有的 FreedomConnector::families 把出站的 address_family 映射为地址族列表:

address_family 绑定的地址族
ipv4_only V4
ipv6_only V6
auto、prefer_ipv4、prefer_ipv6 V4 和 V6

绑定失败的地址族会被略过,并记录一行 debug 日志 udp: no V6 socket: <error>(或 V4)。只有全部都绑定失败时拨号才会失败:AddrNotAvailable udp: no usable local socket in any requested family。不带任何地址族调用 bind_dual 则会以 udp: no address family requested 失败;families 从不返回空列表,所以 freedom 不会遇到它。

DualStackUdp 本身也是一个 DatagramLink,但这个实现只向 IP 目的地发送,对域名以 Unsupported udp: a plain dual-stack link cannot resolve a domain 拒绝,因为如何解析名称是应用层的策略。所以 freedom 把它包进 ResolvingUdp,后者应用出站的解析器和 address_family,并通过 socket 自己的 poll_send_to(SocketAddr) 发送:

  • 发往 IP 的包从该 IP 所属地址族的 socket 发出。
  • 发往域名的包,发往出站解析器按其 address_family 为该域名返回的地址。
  • 发往无法解析的名称的包被丢弃:poll_send_to 报告它已发送(Ok(buf.len())),并以 debug 级别记录 freedom: dropping a datagram to an unresolvable <destination>,目的地以其 Debug 形式输出。
  • poll_recv_from 从有数据报的那个 socket 读取,先检查 v4 再检查 v6,并以 Destination::udp(from) 返回来源。

发往一个没有对应地址族 socket 的地址,会以 AddrNotAvailable udp: no local socket in the family of <peer> 失败,例如 ipv4_only 的 freedom 上的 IPv6 地址。拨号器本身见拨号器与 socket 策略。

两者都自己持有传输,而不是经由 TransportConnector 拨号;两者都在第一个流到来时才建立传输,而不是在构建时。所以两者的构建都不会失败,--test 也从不联系它们的服务器。

Hysteria 2。 build_outbound 这样填写 Hy2Config:

supervisor/src/build/outbound.rs
Hy2Config {
server: udp(&hy2.server), // 服务器是一个 UDP 端点
server_name: hy2.server_name.clone(),
password: hy2.password.expose().clone(),
verify: hy2.verify.clone(),
obfs: hy2.obfs.as_ref().map(obfs),
max_concurrent_streams: hy2.max_concurrent_streams,
socket: socket.clone(),
}

然后构建 Hy2Connector::with_address_family(config, hy2.address_family).with_resolver(dns.servers.clone())。第一个流到来时,connector 的槽拨号建立连接:它用 servers 按 address_family 解析服务器,并按顺序尝试每个地址。每次尝试都用 UdpDialer::new(config.socket).bind(family of that address) 绑定自己的 UDP socket(即按该地址的地址族绑定),设置了 obfs 时把它包进 SalamanderSocket,然后在上面运行 QUIC。名称解析不出可用地址时以 AddrNotAvailable hysteria2: the server name resolved to no usable address 失败,服务器的所有地址都无应答时以 ConnectionRefused hysteria2: no address answered (<addr>: <error>; …) 失败。同一出站后来的流都使用槽中保存的、已认证的连接:TCP 流作为代理 stream(OutboundStream::Hy2),UDP 流作为基于 QUIC 数据报的 UDP 会话。

VerifyMode::CustomCa 的 CA 证书包由 client_config(protocols/src/hysteria/connection.rs)在每次拨号建立连接时解析,而不是在构建出站时。它的错误都是 InvalidInput:hysteria2: could not read the CA file: <error>、hysteria2: the CA file holds an unusable certificate: <error> 和 hysteria2: the CA file contains no certificates。

UDP 流在两种情况下以 Unsupported 被拒绝:服务器的认证应答表示不中继 UDP 时为 hysteria2: the server does not relay UDP,对端从未声明支持 QUIC 数据报时为 hysteria2: the server did not offer QUIC datagrams。任一情况下,Hy2Connector::dial 都会记录一条警告 hysteria2: <error>; datagrams routed to this outbound are dropped(target 为 etemenanki_protocols::hysteria::connector),否则把 UDP 路由到这里的运维者什么也看不到。背后的标志 udp_warned 是一个由 connector 的所有克隆共享的 Arc<AtomicBool>,所以这条警告每个出站版本只记录一次,而不是每个数据报一次。

Hy2DatagramLink 发送每个数据报时把目标写成 authority 形式(format_authority)。接收方向上它重组分片(Defragger),跳过地址无法解析或端口为 0 的回复,并在投递回复的 channel 关闭后以 BrokenPipe hysteria2: the connection behind this association is gone 失败。客户端见 Hysteria 2:协议与客户端。

WireGuard。 build_outbound 这样填写 WgConfig:

supervisor/src/build/outbound.rs
WgConfig {
private_key: *wg.private_key.expose(),
peer_public_key: wg.peer_public_key,
preshared_key: wg.preshared_key.as_ref().map(|k| *k.expose()),
endpoint: udp(&wg.endpoint),
endpoint_resolution: Some((servers.clone(), wg.endpoint_address_family)),
local_addrs: wg.addresses.clone(),
mtu: wg.mtu,
persistent_keepalive: wg.keepalive,
reserved: wg.reserved,
socket: socket.clone(),
}

然后构建 WgConnector::with_address_family(config, wg.address_family).with_resolver(dns.destinations.clone())。这个分支上的注释说明了分工:端点和其他代理服务器一样,所以由 servers 解析;隧道内访问的是本进程自己解析的目的地,所以由 destinations 解析。第一个流到来时,acquire_device 启动设备。local_addrs 为空时,WgDevice::start 以 InvalidInput wireguard: no tunnel-local addresses configured 失败。否则它用 servers 按 endpoint_address_family 解析端点并取第一个地址(没有地址时为 NotFound wireguard: endpoint domain did not resolve),用 UdpDialer::new(config.socket) 在策略下按该地址的地址族绑定一个 UDP socket,把它 connect 到端点,并 spawn 隧道的驱动任务。它不执行 WireGuard 握手,所以即使对端不可达也会成功。

隧道内的目的地由 resolve_candidates("wireguard", …)(protocols/src/helpers/address_family.rs)用 destinations 解析,并按 address_family 和隧道本地地址的地址族(FamilySupport::from_addrs)过滤,因为用户态网络栈没有可查询的路由表。解析不出任何结果的名称得到 NotFound wireguard: destination did not resolve。没有地址剩下时,错误是 AddrNotAvailable wireguard: no usable <address_family> destination address for <host>:<port>;如果是隧道本地地址排除了候选地址,会追加 (local address supports IPv4 only) 或 (local address supports IPv6 only);解析器的错误原样透传。TCP 流按顺序尝试候选地址,隧道内的每次连接受 TCP_CONNECT_ATTEMPT_TIMEOUT(10 秒)限制;全部失败时错误是 TimedOut wireguard: tunnel TCP connect failed for all resolved addresses (<ip>: <error>; …)。隧道见 WireGuard。

连接槽与隧道槽。 两个 connector 都把建立起来的东西保存在一个由所有克隆共享的槽中,并在它失效后重新建立:

Hysteria 2(protocols/src/hysteria/slot.rs) WireGuard(protocols/src/wireguard/slot.rs)
槽 ConnSlot,位于 Arc<parking_lot::Mutex<_>> 之后,状态为 Idle、Connecting 或 Ready DeviceSlot,位于 Arc<tokio::sync::Mutex<_>> 之后,持有 Option<Arc<WgDevice>>
建立 start_connect spawn 一个分离的任务运行 Hy2Conn::connect,并把结果写入槽。期间到达的拨号共享它的结果(Shared future),所以一个出站永远不会同时打开两条连接。connect 期间不持有任何锁。 槽中没有存活的设备时,acquire_device 运行 WgDevice::start,并在返回前把设备存入槽中。
发现失效 Hy2Conn::is_alive 为 false:以 warn 级别记录为 hysteria2: connection closed, reconnecting 驱动任务已结束(WgDevice::is_alive):以 warn 级别记录为 wireguard: tunnel driver stopped, rebuilding
一次失败的尝试 以 warn 级别记录为 hysteria2: connect failed: <error> 以 warn 级别记录为 wireguard: tunnel start failed: <error>
失败后的退避 RECONNECT_BACKOFF_BASE(1 秒)按连续失败次数翻倍,上限为 RECONNECT_BACKOFF_MAX(30 秒):2、4、8、16 秒,然后 30 秒 REBUILD_BACKOFF_BASE(1 秒)和 REBUILD_BACKOFF_MAX(30 秒),时间表相同
短命的连接或隧道 建立后 MIN_HEALTHY_LIFETIME(10 秒)内就失效的算作一次失败;运行得更久的会重置计数并立即重连 相同,使用它自己的 MIN_HEALTHY_LIFETIME(10 秒)
退避期间的拨号 以 BrokenPipe hysteria2: connection is down, waiting before the next attempt 失败 以 BrokenPipe wireguard: tunnel is down, waiting before the next attempt 失败

之所以有短命规则,是因为认证成功后立即被关闭的连接(服务器已达用户上限),或者 start 对着一个失效对端也能成功的隧道,否则会被每个到来的流重新拨号,而退避永远累积不起来。如果 Hysteria 2 的 connect 任务没有发送结果就结束了,等待方会得到 hysteria2: the connect task disappeared。

被复用的出站保留它已建立的东西。an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply 验证了这一点:在复用该出站的一次应用之后打开的流,使用的是第一个流打开的 QUIC 连接;而变化了的 Hysteria 2 出站是一个新版本,会拨号建立自己的连接。

Outbound::Dns(Service) 是 supervisor 自己的 DNS 服务;当 DNS spec 让 supervisor 负责应答时,被拦截的 53 端口流会到达这里。build_outbound 不构建它:prepare 把 Dns::service 包装为 Target::outbound(internal_id(epoch + 1), Outbound::Dns(service))。这个内部 id 的 tag 为空,而任何 spec 的 tag 都不能为空;它只随 DNS spec 重建。服务如何应答见名称解析与 DNS 服务。

supervisor/src/topology/outbound/udp_fanout.rs
pub const MAX_SUBS: usize = 64;
struct Sub {
/// The outbound version the link is open on.
id: OutboundId,
/// The target the route named: the same version, or the balancer that picked it.
via: OutboundId,
link: MeteredDatagram<Guarded<OutboundDatagram>>,
}
pub struct FanOutLink {
plane: PlaneCell,
ctx: FlowContext,
/// The association's identity: every packet routes as this flow, toward its own target.
flow: Flow,
scope: FlowScope,
/// Most recently sent to at the back.
subs: Vec<Sub>,
/// Where the last receive stopped, so every sub gets its turn.
next: usize,
// 私有的打开状态
/// The receiver, parked while no sub can deliver.
recv_waker: Option<Waker>,
}
impl FanOutLink {
pub(crate) fn new(plane: PlaneCell, ctx: FlowContext, flow: Flow, scope: FlowScope) -> Self;
}
impl DatagramLink for FanOutLink {
type Addr = Destination;
// poll_send_to, poll_recv_from
}
supervisor/src/connector.rs
#[derive(Clone)]
pub(crate) struct FlowScope {
pub tracker: Tracker,
pub session: Option<SessionId>,
/// The session's wire, when the sub-links' payload counts toward it.
pub charge: Option<Arc<Wire>>,
}

每个数据报都带有自己的地址,所以 UDP 关联没有单一的目的地可供路由。把它作为整体只路由一次,会让一个出站承载所有对端的流量,UDP 在第一个包之后就会绕过所有路由规则。因此 AppConnector::connect 根本不为 UDP 流拨号:它返回一个 FanOutLink,其中包含 plane cell、流上下文、流本身和一个 FlowScope(tracker、会话 id,以及 connector 以 charging_datagrams 构建时会话的 wire)。服务端核心立即看到 Connected,之后每个包各自路由。见深入 UDP fan-out。

supervisor/src/topology/balancer.rs
pub const DEFAULT_PROBE_INTERVAL: Duration = Duration::from_secs(30);
pub const DEFAULT_PROBE_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Strategy {
/// The first healthy member in configured order. Order is the priority.
Failover,
/// Each healthy member in turn.
RoundRobin,
}
impl Strategy {
pub fn parse(s: &str) -> io::Result<Self>;
}
pub struct Member {
pub tag: CompactString,
pub target: Arc<Target>,
pub probe: Destination,
healthy: AtomicBool,
}
impl Member {
pub fn new(tag: CompactString, target: Arc<Target>, probe: Destination) -> Self;
pub fn is_healthy(&self) -> bool;
pub fn set_healthy(&self, up: bool) -> bool;
}
pub struct Balancer {
members: Vec<Arc<Member>>,
strategy: Strategy,
next: AtomicUsize,
}
impl Balancer {
pub fn new(members: Vec<Arc<Member>>, strategy: Strategy) -> io::Result<Self>;
pub fn select(&self) -> Arc<Target>;
pub fn select_preferring(&self, prefer: Option<&OutboundId>) -> Arc<Target>;
pub fn spawn_probe(
self: &Arc<Self>,
tracker: &TaskTracker,
token: CancellationToken,
interval: Duration,
timeout: Duration,
resolver: Resolver,
dialer: TcpDialer,
);
}

负载均衡器不是 Outbound。它以 TargetKind::Balancer(Arc<Balancer>) 类型的 Target 存在于 plane 中,使用自己的 tag;这个 tag 与出站共享同一个命名空间,所以路由指向负载均衡器的方式与指向出站完全一样。成员的 target 就是同一次应用中 plane 为该出站 tag 持有的那个 Arc<Target>,所以成员仍可以通过自己的 tag 单独使用,而经过负载均衡器的流登记在成员的版本之下。spec 一侧是 supervisor/src/topology/spec_plan/route.rs 中的 BalancerSpec { tag, members, strategy, probe_interval, probe_timeout }。见深入负载均衡器。

sequenceDiagram
  participant RT as 服务端运行时
  participant A as AppConnector
  participant T as Target
  participant O as Outbound
  participant C as ProxyClient
  RT->>A: connect(flow)
  A->>A: 在会话上准入,加载 plane,路由
  A->>T: resolve(None),负载均衡器此时选出成员
  A->>T: connect_stream(flow)
  T->>O: connect_stream(flow)
  O->>C: 在锁内 connect(flow)
  C-->>O: ProxyClientConnecting
  O-->>RT: 装箱的 future,经由 Target 和 AppConnector 返回
  RT->>RT: poll:拨号、codec 握手、flush
  alt 已打开
    RT->>RT: proxy_stream,然后 Guarded,然后 Metered
  else 任何错误
    RT->>RT: 向核心发送 Event ConnectFailed
  end

此后服务端运行时直接 poll 这个 stream。对代理客户端而言,这意味着在连接自己的任务里 poll 一个 ProxyClientRuntime:入站和上游之间没有任何任务或 channel。图中路由的那一半见 plane:为每个流选路。

sequenceDiagram
  participant C as 服务端核心
  participant F as FanOutLink
  participant P as Plane
  participant T as Target
  participant S as 子 link
  C->>F: 带 UDP 流的 Effect Open,不拨号
  F-->>C: 立即返回 Event Connected
  C->>F: poll_send_to(packet, peer)
  F->>P: load,为指向 peer 的流 route
  P-->>F: 目标和规则
  F->>T: resolve(prefer)
  alt 已有使用这个 id 的子 link
    F->>S: poll_send_to
  else 还没有子 link
    F->>T: 为指向 peer 的流 connect_datagram
    T-->>F: Guarded link
    F->>F: 作为新流计量,加入表中
    F->>S: poll_send_to
  end
  C->>F: poll_recv_from
  F->>S: 从 next 开始 poll 每个子 link
  S-->>F: 负载和来源
  F-->>C: 来自该来源的 Event Datagram

poll_send_to(cx, buf, to) 以同样的方式开始处理每个包:

supervisor/src/topology/outbound/udp_fanout.rs
let flow = self.flow.toward(to.clone());
let plane = self.plane.load();
let routed = plane.route(&flow, &self.ctx);
let rule = routed.rule;
let routed = routed.target.clone();
  • Flow::toward 把关联的用户和来源复制到一个指向本包目的地的流中,并清除 sniffed。因此逐包路由看到的是包自己的地址(IP 或域名,取决于客户端发送的是什么)、入站 tag、来源和 UDP 网络。域名或 geosite 规则只有在客户端按名称寻址某个 UDP 包时才会匹配它。
  • 每个包都会加载一次 plane,所以路由变化或新的出站版本会作用到活跃关联的下一个包(a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow)。
  • 存在 DNS 服务时,Plane::route 为它拦截 53 端口,所以 UDP DNS 查询会成为内部 DNS 目标上的一个子 link。
  • 每个包的工作量是:Flow::toward(克隆目的地和用户的 Arc)、一次 ArcSwap 加载、路由模型中描述的路由表首条匹配遍历,以及对子 link 表(最多 MAX_SUBS 项)的线性扫描,用于找到首选的和匹配的子 link;发送成功后该子 link 被移到末尾。
supervisor/src/topology/outbound/udp_fanout.rs
let prefer = self
.subs
.iter()
.rev()
.find(|s| s.via == *routed.id())
.map(|s| &s.id);
let target = routed.resolve(prefer);
let id = target.id();

prefer 是由同一路由目标引出的、最近使用的子 link 的版本。Target::resolve 对出站目标原样返回,对负载均衡器则调用 select_preferring(prefer)。这段代码的注释解释了为什么在这里而不是在打开时解析:子 link 以承载该包的版本为键,所以成员不可用时,关联在下一个包就会转到另一个成员,而不是停留在登记于负载均衡器名下的子 link 上。

  • 在 round_robin 下,首选成员只要健康就会保留,所以关联不会逐包在成员之间跳动(round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down)。
  • 在 failover 下,首选被忽略:每次都选第一个健康成员,因为回到恢复了的更高优先级成员正是它的目的(failover_ignores_the_preference)。
  • 首选是某个成员的某一个版本。该成员被重建后,这个首选不再对应任何成员,选择就像没有首选一样进行(a_preference_for_another_version_is_no_preference)。

子 link 通过 id 查找,即解析出的出站的 OutboundId(tag 和版本)。键从不是目标的内存地址:出站的生命周期长于单次应用,一个 tag 可能同时有旧版本和新版本,而释放的内存可能被它的后继复用(supervisor/src/entity/id.rs)。由此:

  • 路由到同一出站版本的两个对端共享一个子 link。对 freedom、SOCKS、Trojan、WireGuard 和 Hysteria 2 而言,这就是一个 socket 或一个代理关联承载多个对端。
  • 通过自己的 tag 直接访问的成员,与经由负载均衡器访问的同一成员,id 相同,所以共享一个子 link。
  • 一个 tag 的新版本是一个新键,所以路由到它的下一个包会打开自己的子 link。旧版本的排空见 plane:为每个流选路。
  • blackhole 也是一个子 link:一个吞掉包的 BlackholeLink。被屏蔽对端的包就是这样被丢弃的,不影响关联的其余部分。

via 记录引出这个子 link 的路由目标:出站本身,或者选中它的负载均衡器。每次发送成功都会把它改写为路由该包的目标,所以之后路由到那里的包也会首选这个子 link。

flowchart TB
  A["在当前 plane 上路由这个包"] --> B["解析负载均衡器,首选正在使用的成员"]
  B --> C{"有使用这个 id 的子 link?"}
  C -- "是" --> D["sub.poll_send_to"]
  D -- "Ready Ok" --> E["移到末尾,via = 路由目标"]
  D -- "Ready Err" --> F["移除它,丢弃这个包,返回 Ok"]
  C -- "否" --> G["用 Target::connect_datagram 打开"]
  G -- "已打开" --> H["计量,加入表中,在上面发送这个包"]
  G -- "失败" --> I["记录日志,丢弃这个包,返回 Ok"]

poll_send_to 实现的规则:

  • 最近最少使用顺序。 subs 按最后一次成功发送排序,最近的在末尾。新打开的子 link 在首次发送之前就加入末尾。只有成功的发送才会移动子 link。
  • 打开失败会丢弃它的包。 为该包的 id 打开失败时,poll_send_to 丢弃这个包并返回 Ready(Ok(buf.len())),所以核心看不到错误。失败以 debug 级别记录为 udp fan-out: opening <id> failed: <error>,例如 udp fan-out: opening web@v1 failed: http carries no datagrams。
  • 发送失败会丢弃子 link 和它的包。 子 link 的 poll_send_to 返回错误时,移除该子 link,把 next 重置为 0,以 debug 级别记录 udp fan-out: sending on <id> failed: <error>,并返回 Ready(Ok(buf.len())),就像有损链路丢了一个包。之后路由到那里的包会打开一个新的子 link。子 link 的出站版本被排空策略关闭时(Guarded 以 the outbound this flow was opened on was drained 失败),或者它的流通过 tracker 被强制关闭时,子 link 也是这样结束的。
  • 没有错误会到达核心。 FanOutLink::poll_send_to 从不返回错误,所以服务端运行时从不为 fan-out 关联报告 SendFailed。

没有子 link 使用该包的 id 时,fan-out 用 Box::pin(target.connect_datagram(flow.clone())) 打开一个,并随之保存 id、via,以及该子 link 被跟踪时使用的 FlowMeta。FlowMeta 是 FlowMeta::of(&flow, session, &ctx.inbound_tag, ctx.source, id, rule, plane.epoch()):指向打开该子 link 的那个包的流、路由那个包的规则,以及它被路由时所在 plane 的 epoch。由于目标已经解析过,Target::connect_datagram 不会再次询问负载均衡器。

poll_opening 是子 link 加入表的唯一位置。打开成功时,它:

  1. 用 scope.tracker.meter_datagram(meta, link, scope.charge.clone()) 包装 link;
  2. 把新的 Sub 推到 subs 末尾;
  3. 如果 subs.len() > MAX_SUBS,移除索引 0,即最久未被发送(或打开)的子 link,并用 saturating_sub(1) 递减 next。丢弃被移除的子 link 会结束它的被跟踪流并丢弃它的 link;
  4. 唤醒 recv_waker,让在空表上等待的接收方 poll 新的子 link。

打开失败时,它以 debug 级别记录失败。表中打开的子 link 永远不超过 MAX_SUBS(64)个。

poll_recv_from 从 next 开始循环,把每个子 link 各 poll 一次:

  • Ready(Ok(from)):把 next 设为下一个索引并返回 from。下一次接收从这个子 link 之后开始,所以一个繁忙的对端不会饿死其他对端。
  • Ready(Err(e)):子 link 已结束,例如因为 SOCKS 控制 stream 关闭、代理 stream 失败、它的版本以 Close 被排空,或它的流被强制关闭。它被移除;如果 next 越过了末尾,则重置为 0;错误以 debug 级别记录为 udp fan-out: sub-link <id> ended: <error>;任务唤醒自己,以便再次 poll 剩下的子 link。
  • 全部返回 Pending:把 cx.waker() 存入 recv_waker 并返回 Pending。

FanOutLink::poll_recv_from 从不返回错误。服务端运行时把接收错误视为出站的致命错误,而在这里出站就是整个关联:一个出问题的对端不能让它结束。因此 fan-out 从不因为某个出站失败而结束关联;关联在它的连接结束时结束,例如入站一侧结束了它,或者 supervisor 关闭了它所属的会话(见用户与会话)。

每个子 link 都是一个独立的被跟踪流,即 MeteredDatagram<Guarded<OutboundDatagram>>:

  • 它登记在关联的会话、解析出的出站版本和路由其第一个包的规则之下,所以在 sessions() 和流表中,一个关联用到几个出站就显示几个流(a_udp_association_counts_each_sub_link_and_charges_its_session:发往两个出站的三个包得到两个流)。
  • 它统计在它上面发送和接收的负载字节,并执行会话的限速。
  • 设置了 FlowScope::charge 时,它的负载也计入会话的线路字节(wire bytes)。serve_connection(supervisor/src/serve.rs)用 charging_datagrams() 构建 SOCKS 入站的 connector,因为 SOCKS 数据报从不经过会话自己的 TCP stream。其他入站的数据报已经经过了其会话所统计的东西:代理 stream,或 Hysteria 2 会话的 QUIC 连接。TUN 流没有会话。
  • 通过 tracker 强制关闭它,只会结束这一个子 link,和任何失败的子 link 结束的方式一样:关联继续运行,之后路由到那里的包会打开一个新流。

计量、限速和强制关闭见流跟踪、统计与限速;会话线路字节见按用户的用量计费。

Section titled “各子 link 如何处理每个包的地址”
出站 子 link 是否遵从每个包的目的地
freedom 是:ResolvingUdp 把每个包发往它自己的地址。
blackhole 不适用:包被吞掉。
socks 是:每个包都带有自己的 SOCKS5 UDP 头。
trojan 是:每个 Trojan UDP 包都带有自己的地址。请求头指明 UDP 命令(0x03)和一个 0.0.0.0:0 占位值。
vless、vmess 否:请求头只指明一个目标,即打开该子 link 的那个包的目的地;seal_to 忽略每个包的地址,open_from 把每个回复都归到同一个目标。
wireguard 是:每个数据报发往隧道内它自己的地址。
hysteria2 是:每个数据报带有自己的目标,域名以名称形式交给服务器解析。
DNS 服务 由 DNS 服务在进程内应答。
http、shadowsocks(传统与 2022 两类) 不支持 UDP:打开以 Unsupported 失败,包被丢弃。

validate(supervisor/src/build/validate.rs)在构建任何东西之前检查每个负载均衡器,所以 --test 能发现以下全部问题。错误文本是本文核对版本的 etemenanki-app 在 configuration invalid: 之后打印的内容:

规则 错误
负载均衡器的 tag 不为空 balancer : a balancer tag must not be empty
没有两个负载均衡器共用一个 tag duplicate balancer tag <tag>
负载均衡器的 tag 不同时是某个出站的 tag duplicate balancer tag <tag>
负载均衡器至少有一个成员 balancer <tag> has no members
每个成员都指向一个出站 balancer <tag> references unknown outbound <member>
每个成员都有一个 TCP 探测可以到达的上游(probe_target) balancer <tag>: outbound <member> has no upstream a TCP health probe can reach
app 配置设置了 strategy 时,它必须恰好是 failover 或 round_robin balancer <tag>: unknown balancer strategy "<value>" (expected "failover" or "round_robin")

probe_target 对七个代理客户端(SOCKS、HTTP、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022)返回 upstream.server,对其余的返回 None。freedom 和 blackhole 没有上游。Hysteria 2 和 WireGuard 监听 UDP,所以 TCP 连接会让它们永远被标记为不可用。成员只在 spec 的出站中查找,所以负载均衡器不能是另一个负载均衡器的成员;这种情况得到的是“unknown outbound”错误。

strategy 的错误文本来自 Strategy::parse:etemenanki-app 在 lowering(降为 spec)时由 app/src/lower.rs → lower_balancer 调用它,并加上前缀 balancer <tag>: 。同一个 lowering 把未设置的 strategy 变成 Strategy::Failover,把未设置的 probe_interval 和 probe_timeout(整数秒)变成 DEFAULT_PROBE_INTERVAL 和 DEFAULT_PROBE_TIMEOUT。spec 本身携带 Strategy 和两个 Duration。app 的配置见etemenanki-app:从 TOML 到 spec。

新的 plane 发布后,commit 为本次应用构建的每个负载均衡器启动探测:

supervisor/src/supervisor.rs
let token = self.root.child_token();
balancer.spawn_probe(
&self.shared.tracker,
token.clone(),
interval,
timeout,
dns.servers.clone(),
TcpDialer::new(self.socket.clone()),
);
self.probes.insert(tag, token);

这些任务运行在 supervisor 的 TaskTracker 上,归属于 supervisor 根 token 的一个子 token,Actor::probes 按负载均衡器 tag 保存这个子 token。它们用 servers 解析,即成员上游拨号时所用的解析器,并在 socket 策略下连接。

spawn_probe 为每个成员 spawn 一个任务:

stateDiagram-v2
  [*] --> Probing
  Probing --> Recording: 已连接、失败或超时
  Probing --> [*]: token 被取消,放弃探测
  Recording --> Sleeping: set_healthy,状态变化时记录日志
  Sleeping --> Probing: 间隔已到
  Sleeping --> [*]: token 被取消
  • probe(dest, timeout, resolver, dialer) 以 AddressFamilyStrategy::Auto 解析 member.probe,并依次对每个地址尝试 dialer.connect。任一地址接受连接即为可用;连接随即丢弃。tokio::time::timeout 以 timeout 限制整个探测,包括解析;每次连接尝试还受拨号器的 DEFAULT_CONNECT_TIMEOUT 限制。解析错误、没有地址或超时都算作不可用。
  • 探测只是一次 TCP 连接,仅此而已。 它既不使用成员的传输层,也不使用它的凭据或 address_family:它回答的是“上游是否可达”,而这正是负载均衡器要绕开的那类故障。更复杂的探测需要凭据,也会让探测变成代理客户端的第二套实现。
  • set_healthy 以 Ordering::Relaxed 交换成员的 AtomicBool 并返回旧值。状态变化以 info 级别记录为 balancer member <tag> is now up 或 balancer member <tag> is now down。
  • 成员初始为健康。 Member::new 把 healthy 设为 true,第一次探测在任务启动后立即运行。从未探测过的上游算作健康,因为如果每个成员初始都不可用,每次启动后的整整一个间隔内都会丢弃所有流量。
  • 取消会放弃正在进行的探测。 探测和休眠都在 tokio::select! 中与 token.cancelled() 竞争,所以取消会立即结束任务,而不必等探测完成。

端到端测试 traffic_moves_off_a_member_that_stops_answering 以 probe_interval = 1 运行两个 app 实例。第一个成员接受 TCP 但不是代理,所以它的探测结果是健康的,而经过它的流会失败;一旦它停止接受连接,探测就把它标记为不可用,流转到可以工作的成员上。

supervisor/src/topology/balancer.rs
pub fn select(&self) -> Arc<Target> {
let mut healthy = self.members.iter().filter(|m| m.is_healthy());
let chosen = match self.strategy {
Strategy::Failover => healthy.next(),
Strategy::RoundRobin => match healthy.clone().count() {
0 => None,
n => healthy.nth(self.next.fetch_add(1, Ordering::Relaxed) % n),
},
};
chosen.unwrap_or(&self.members[0]).target.clone()
}

当策略为 RoundRobin、prefer 为 Some,且 target.id() 与之相等的成员是健康的时,select_preferring(prefer) 返回该首选成员的目标。否则它调用 select。

调用方:

  • TCP 流:AppConnector::connect 调用一次 Target::resolve(None),流登记在选中的成员之下并受其守护。Target::connect_stream 会再解析一次,对出站目标而言结果就是目标本身。
  • UDP 包:FanOutLink::poll_send_to 为每个包调用 Target::resolve(prefer),如上文所述。

选择从不在路由时进行,所以成员不可用影响的是下一个流或下一个包,而不是下一次应用。

flowchart TB
  P{"round_robin,且首选成员健康?"} -- "是" --> K["保留首选成员"]
  P -- "否" --> A["healthy 标志已设置的成员"]
  A --> B{"有健康的成员吗?"}
  B -- "否" --> F["第一个配置的成员"]
  B -- "是" --> C{"策略"}
  C -- failover --> D["按配置顺序的第一个健康成员"]
  C -- round_robin --> E["索引为 next.fetch_add(1) mod count 的健康成员"]
情况 failover round_robin
多个成员健康 按配置顺序的第一个健康成员。更高优先级的成员恢复后,新流回到它上面。 轮流使用每个健康成员;UDP 关联在成员保持健康期间一直使用它。
部分成员不可用 跳过。 跳过。计数器持续递增,所以健康集合变化时轮转顺序会偏移。
所有成员都不可用 第一个配置的成员。 第一个配置的成员。

退回第一个成员而不是让流失败,是有意为之:发往一个可能已经恢复的上游的流,好过一个注定被丢弃的流。另一种做法是:从每个成员最近一次探测都失败的那一刻起,到其中一个再次成功为止,负载均衡器丢弃每一个流,而在上游恢复之后,这种状态最多还会持续整整一个 probe_interval。选择是无锁的:对成员健康标志做 relaxed 原子加载,在 round_robin 下若有任何成员健康,再加一次 relaxed fetch_add。failover 在第一个健康成员处停止加载;round_robin 最多遍历两遍,一遍统计健康成员数(healthy.clone().count()),一遍到达被选中的成员(healthy.nth(k))。

负载均衡从不重试。分发会消耗掉流,所以在选中成员上失败的拨号会以 ConnectFailed(TCP)或被丢弃的包(UDP)的形式到达核心。能否绕开失效的上游完全取决于探测是否已经发现它,这也是为什么健康状态在带外建立,与 Xray 的 observatory 相同。

经过负载均衡器的流在成员的 Target 上打开,并由成员版本的关闭 token 守护,而不是负载均衡器的。a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers 验证了这一点:关闭负载均衡器自己的目标不影响这个 stream,关闭成员的目标则会中止它。计划只为出站版本产生 Drain 步骤,所以负载均衡器本身永远不会被排空。

负载均衡器持有其成员的目标,所以任何一个成员重建时,plan 都会重建它:

supervisor/src/topology/spec_plan/plan.rs
let members_kept = !balancer.members.iter().any(|m| rebuilt.contains(m));
match running_version {
Some(version) if prev == Some(balancer) && members_kept => {
// 在运行中的版本上 Step::Reuse(Resource::Balancer(tag))
}
_ => {
// 在下一个版本上 Step::Build(Resource::Balancer(tag))
}
}
变化 负载均衡器会怎样
负载均衡器的 spec 和所有成员都没变 Step::Reuse。保留运行中的 Arc<Target>,也就保留了其中的 Balancer:已知的健康状态、轮询计数器和探测任务,commit 保留这些任务的 token(probes.retain 保留计划复用的每个 tag)。
负载均衡器的 spec 变了,或任一成员被重建 下一个版本上的 Step::Build。新的 Balancer 以所有成员健康开始,并在一个新的子 token 下获得新的探测任务;旧 token 在同一次 commit 中被取消。
DNS spec 变了 每个成员都被重建(每个可探测的出站都使用 DNS),所以每个负载均衡器都被重建并重新获知健康状态。
负载均衡器被移除 plane 在去掉它之后重新发布,它的探测 token 被取消。它的版本被保留,所以以同一 tag 重新加回的负载均衡器得到版本 2(removing_a_balancer_alone_republishes_the_plane)。

a_balancer_is_rebuilt_exactly_when_one_of_its_members_is 验证了第一行和第二行中关于成员的那一半;负载均衡器自身 spec 的变化没有专门的测试。探测的启动和停止发生在 commit 中发布 plane 的过程里,而每次负载均衡器的构建或移除都会触发 plane 发布。关闭时,Actor::shutdown 先取消每个探测 token,再关闭 task tracker 并等待它。

socket 策略是一个 SocketOptions,通过 SupervisorBuilder::socket_options 在 builder 上设置一次。它在 supervisor 的整个生命周期内有效;应用无法改变它。它的文档注释列出了覆盖范围:出站或负载均衡器探测的每次拨号、直连流的每个 UDP socket、SOCKS UDP 中继、Hysteria 2 或 WireGuard 隧道,以及解析器从本机发出的每个查询。移动端 VPN 正是在这里把代理自己的流量排除在它的隧道之外。

Socket 由谁打开 策略如何作用到它
freedom TCP 连接 TcpDialer::connect_any build_outbound 中的 Dialer::new(socket)
freedom UDP socket,每个地址族一个 UdpDialer::bind_dual 同一个 Dialer
代理上游 TCP:HTTP、SOCKS(CONNECT 和关联的控制 stream)、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022 TransportConnector::dial transport 中的 Dialer::new(socket)
SOCKS UDP 中继 socket SocksUdpLink::associate 的 bind 闭包 build_outbound 中的 UdpDialer::new(socket)
Hysteria 2 QUIC socket protocols/src/hysteria/connection.rs 中的 bind_socket Hy2Config.socket
WireGuard 隧道 socket WgDevice::start WgConfig.socket
负载均衡器探测连接 probe commit 中的 TcpDialer::new(self.socket)
发往配置的 DNS 服务器的查询 解析器的上游 build_dns 中的 ResolverOptions.socket(名称解析与 DNS 服务)
经由代理发出的 DNS 查询 经由路由表的 RoutedDialer 承载它们的那个出站的 socket
系统解析器(getaddrinfo) C 库 覆盖不到:它在策略无法触及的地方打开自己的 socket

SocketOptions::apply 在创建 socket 之后、绑定或连接之前运行,这是大多数选项唯一能设置的时机。它先设置接口,再设置 mark(SO_MARK),然后运行 hook。接着 TcpDialer::socket 在策略有 tcp_keepalive 空闲时间时设置它,并在 bind_address 属于该 socket 的地址族时绑定它。在代理上游上,这个 keepalive 不会保留下来:连接之后,TransportConnector::dial 调用 set_keepalive,把它替换为固定的 120 秒、30 秒、3 个探测包的时间表(见代理客户端之下的传输层)。freedom 的 TCP 连接只有策略设置的 keepalive。UdpDialer::bind 在 apply 之前把 IPv6 socket 设为 v6-only,之后把 socket 设为非阻塞,然后在临时端口上绑定同一地址族的 bind_address 或通配地址。平台不支持某个被要求的选项时,以 Unsupported 失败,而不是悄悄忽略。完整的策略见拨号器与 socket 策略。

返回错误的 hook 会让这个 socket 失败,拨号也随之失败:没有任何数据会从策略不接受的 socket 发出。对 freedom TCP,hook 的错误会出现在 failed to connect to any address (…) 中,SOCKS 客户端看到的是一个失败回复(a_refusing_hook_fails_the_dial)。hook 拒绝其 socket 的 UDP 地址族会被 bind_dual 略过,只有在没有地址族剩下时拨号才失败。移动端绑定(ffi/src/platform.rs → socket_options)安装的 hook 会请求平台保护每个 socket 的文件描述符,平台拒绝时以 PermissionDenied the platform refused to protect the socket 失败。

check 用 SocketOptions::default() 构建它要校验的 supervisor,而构建过程不打开任何 socket,所以 --test 从不运行 hook。

  1. 在 OutboundProtocolSpec(supervisor/src/topology/spec_plan/outbound.rs)中添加一个变体。基于 stream 传输层的代理客户端接受一个 ProxyUpstream { server, transport }。规划器用 == 比较 spec,所以每个影响行为的字段都必须在其中。
  2. 在 validate_outbound(supervisor/src/build/validate.rs)中添加它的规则;有 ProxyUpstream 时还包括 outbound_transport。决定 probe_target 返回什么:TCP 上游让它可以做负载均衡,None 则拒绝它成为成员。
  3. 检查 plan 中 uses_dns 的 match:DNS spec 变化时,除 blackhole 外的每个出站都会重建。不解析任何东西的协议可以在那里与 blackhole 并列。
  4. 添加一个 Outbound 变体(supervisor/src/topology/outbound/mod.rs)及其两个分支。stream 经过 ProxyClient 和 proxy_stream(成为 OutboundStream::Proxy),如果不应装箱,则在 OutboundStream 中获得自己的变体。数据报 link 用 OutboundDatagram::new 装箱;不支持 UDP 的协议返回 no_udp("<name>")。
  5. 对 ProxyClient,添加一个 BUF 常量,它要能容纳该协议解开的最大帧,并且大于 codec 的 STAGING_RESERVE,否则 ProxyClientRuntime::new 会在第一个流上 panic。
  6. 添加 build_outbound 分支。把 socket 策略传给它可能打开的每个 socket(Dialer::new(socket)、UdpDialer::new(socket),或协议配置中的 socket 字段),用 dns.servers 解析上游名称,用 dns.destinations 解析它自己访问的目的地,并在 make 闭包之外一次性派生密钥。
  7. 让前端程序能把它降为 spec:etemenanki-app 的 app/src/lower.rs,以及面板能表达它时 katana 的 lowering(katana:入站与出站的构建)。如果订阅文件可以指定它,就在 app/src/subscribe.rs 的 node_outbound 中添加一个分支;该函数把每个订阅节点变成 [[outbound]] 表会给出的 OutboundConfig,于是节点随后经过同样的 lowering;订阅格式见订阅文件。
  8. 测试它:为每种新的 socket 添加一个 socket_policy.rs 用例;如果它在多个流之间保持连接,添加一个 hot_swap.rs 用例;并为每条新规则添加 validate.rs 用例。
不变量 由谁保证 由谁验证
UDP 关联逐包路由,使用该包到来时的当前 plane AppConnector::connect 为 UDP 流返回 FanOutLink;poll_send_to 为每个包加载 plane 并路由 one_association_routes_each_peer_separately(e2e_udp_route.rs)、a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow(hot_swap.rs)
被屏蔽的对端不影响其关联的其余部分 被屏蔽的包进入 BlackholeLink 子 link;poll_recv_from 从不返回错误 one_association_routes_each_peer_separately(第三次交换)
回复归属于发送它的对端 每个子 link 的 poll_recv_from 返回其来源;fan-out 原样传递 replies_from_several_peers_merge_back_correctly(e2e_udp_route.rs)
UDP 规则看到的是 UDP 网络 逐包的流保留 DialNetwork::Udp network_separates_tcp_from_udp(e2e_route_context.rs)
关联用到的每个出站是一个被跟踪流 poll_opening 用 meter_datagram 计量每个子 link a_udp_association_counts_each_sub_link_and_charges_its_session(tracking.rs)
每个关联最多打开 MAX_SUBS 个子 link 推入后超过 MAX_SUBS 时,poll_opening 移除索引 0 没有直接测试
打开或发送失败只丢弃该包,从不让关联失败 poll_send_to 在两种情况下都返回 Ready(Ok(buf.len())) 没有直接测试
子 link 的键从不混淆 键是 OutboundId(tag 和版本),从不是地址 由构造保证
ProxyClient 的锁从不跨 I/O 持有 ProxyClient::connect 返回拨号 future;它在锁外被 poll 由构造保证
代理拨号只在上游拨号、codec 握手和 flush 都完成后才完成 ProxyClientConnecting poll 运行时,直到 poll_connected 就绪 server_runtime_relays_through_a_client_runtime_in_one_task、upstream_dial_failure_is_connect_failed_not_connected、refused_upstream_handshake_is_connect_failed_not_connected(concepts/tests/client.rs)
SOCKS UDP 回复只接受来自中继的 endpoint(from) != endpoint(self.relay) 时,SocksUdpLink::poll_recv_from 跳过该数据报 udp_link_ignores_datagrams_not_from_the_relay、udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay(protocols/tests/pipeline/socks.rs)
不能承载 UDP 的出站拒绝 UDP connect_datagram 为 http、shadowsocks 和 shadowsocks-2022 返回 Unsupported 没有直接测试
出站 socket 在 socket 策略下打开 build_outbound 中的每个拨号器和每个探测拨号器都从 self.socket 构建 direct_flows_are_dialed_on_hooked_sockets、a_socks_upstream_is_dialed_on_hooked_sockets(socket_policy.rs);Hysteria 2、WireGuard 和探测的 socket 没有覆盖
拒绝的 hook 使拨号失败 SocketOptions::apply 在任何连接之前返回 hook 的错误 a_refusing_hook_fails_the_dial(socket_policy.rs)
未变化的出站在应用前后保留它持有的东西 prepare 复用运行中的 Arc<Target> an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply(hot_swap.rs)
变化了的出站被构建为新版本供新流使用,Close 排空会关闭旧版本上的流 plan 构建下一个版本并为旧版本产生 Drain 步骤;Guarded 在其版本关闭后失败 a_changed_outbound_keeps_its_old_flows_unless_drained_with_close(hot_swap.rs,一个 freedom 出站)
负载均衡的流随其成员排空,而不是随负载均衡器 Target::resolve 在流被守护之前选出成员 a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers(plane.rs)
failover 取第一个健康成员,并回到恢复了的成员 Strategy::Failover 下的 select failover_takes_the_first_healthy_member_in_order(balancer.rs)
round_robin 跳过不可用的成员 select 在健康子集上轮转 round_robin_cycles_only_through_healthy_members
所有成员都不可用的负载均衡器仍然会选出成员 select 退回 members[0] every_member_down_still_selects_rather_than_dropping
round_robin 保留健康的首选成员,成员不可用时放弃它 select_preferring round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down
failover 忽略首选 select_preferring 先检查策略 failover_ignores_the_preference
对另一个版本的首选等于没有首选 select_preferring 比较 OutboundId a_preference_for_another_version_is_no_preference
负载均衡器从不为空 validate(EmptyBalancer)和 Balancer::new a_balancer_needs_a_member(validate.rs)、a_balancer_needs_at_least_one_member(balancer.rs)
策略名称必须精确 Strategy::parse 只接受 failover 和 round_robin strategy_names_are_validated
只有带 TCP 上游的出站可以做负载均衡 probe_target 对 freedom、blackhole、Hysteria 2 和 WireGuard 返回 None a_balancer_member_needs_an_upstream_a_tcp_probe_can_reach(validate.rs)、a_member_with_no_upstream_is_refused(e2e_balancer.rs)
成员是出站,从不是负载均衡器 成员只在 spec 的出站中查找 a_balancer_member_may_not_be_another_balancer(validate.rs)
负载均衡器的 tag 不是出站的 tag validate 把重叠作为重复拒绝 a_balancer_may_not_share_an_outbounds_tag(validate.rs)
负载均衡器保留其健康状态,除非它或某个成员被重建 members_kept 时 plan 复用它 a_balancer_is_rebuilt_exactly_when_one_of_its_members_is(plan.rs)
停止应答的成员不再接收流 探测把它标记为不可用;select 跳过它 traffic_moves_off_a_member_that_stops_answering(e2e_balancer.rs)
探测随负载均衡器停止 每个探测任务都在其负载均衡器的 token 上 select,该 token 在重建、移除和关闭时被取消 没有直接测试
位置 条件 Kind 消息
prepare 下面任何一个构建错误 无(包装) building outbound <tag>@v<n> failed: <error>
transport TLS shape 没有 TLS 材料、WebSocket 没有 host、gRPC 没有 authority(会先被校验拒绝) InvalidInput missing tls、missing ws host、missing grpc authority
transport → ClientConfig::with_verify_mode 无法解析的 CA 证书包 Other OpenSSL 的错误文本
transport → ClientConfig::with_verify_mode 没有证书的 CA 证书包 InvalidInput no certificate in CA PEM bundle
TransportConnector::dial UDP 目的地(supervisor 从不传入) Unsupported a proxy transport carries no datagrams of its own
Balancer::new 没有成员(会先被校验拒绝) InvalidInput a balancer needs at least one outbound,包装为 building balancer <tag> failed: …
Strategy::parse 未知的策略名 InvalidInput unknown balancer strategy "<value>" (expected "failover" or "round_robin")
Target::connect_stream、connect_datagram 本身是负载均衡器的成员(校验使其不可能出现) Other balancer member <id> is itself a balancer
connect_datagram http、shadowsocks、shadowsocks-2022 Unsupported http carries no datagrams、shadowsocks carries no datagrams、shadowsocks-2022 carries no datagrams
connect_stream、connect_datagram freedom、wireguard 或 hysteria2 产生了另一种类型 Other freedom dialed a TCP flow as UDP、freedom dialed a UDP flow as TCP,wireguard 和 hysteria2 同理
proxy_stream、proxy_datagram 客户端产生了另一种类型 Other a TCP flow was dialed as datagrams、a UDP flow was dialed as a stream
ProxyClientRuntime 上游在握手期间关闭 UnexpectedEof upstream closed during the handshake
ProxyClientRuntime 上游关闭时还有未读完的半帧 UnexpectedEof upstream closed inside a frame
ProxyClientRuntime 大于 BUF 的上游帧 InvalidData upstream frame larger than the client runtime's buffer
ProxyClientRuntime 封装后的帧装不进暂存缓冲区的数据报 InvalidInput frame larger than the client runtime's buffer
ProxyClientRuntime codec 拒绝了上游的字节 InvalidData codec 自己的错误
ProxyClientRuntime codec 违反了约定 Other、InvalidInput codec opened an impossible frame、codec consumed more reply than it was given、codec handshake step made no progress、codec sealed an impossible byte count、codec refused a packet that fits its reserve
ProxyClientRuntime 传输层给出的是数据报 socket Unsupported proxy client runtime needs a stream to the upstream, got a datagram socket
ProxyClientRuntime 线路失败之后的任何调用 NotConnected 无
freedom TCP,解析 名称解析不出任何结果 NotFound dial: destination did not resolve
freedom TCP,解析 没有属于允许地址族的地址 AddrNotAvailable dial: no usable <strategy> destination address for <host>:<port>
freedom TCP,连接 某个地址没有及时应答 TimedOut connect to <addr> timed out(被收集起来,不单独返回)
freedom TCP,连接 没有地址能连上 ConnectionRefused failed to connect to any address (<addr>: <error>; …)
freedom UDP 没有地址族能绑定 AddrNotAvailable udp: no usable local socket in any requested family
UdpDialer::bind_dual 没有请求任何地址族(freedom 从不这样调用) AddrNotAvailable udp: no address family requested
作为 DatagramLink 的 DualStackUdp 域名目的地(freedom 把它包进 ResolvingUdp) Unsupported udp: a plain dual-stack link cannot resolve a domain
ResolvingUdp 发送 该地址的地址族没有 socket AddrNotAvailable udp: no local socket in the family of <peer>
ResolvingUdp 发送 名称无法解析 无 包被丢弃,并记录一行 debug 日志
SocksUdpLink::associate 方法被拒绝、账号被拒绝 PermissionDenied auth method not supported、server rejects account
SocksUdpLink::associate 非零的回复状态 ConnectionRefused server rejects request: <status>
SocksUdpLink::associate 中继以域名给出 Unsupported socks: the relay address is a domain
SocksUdpLink::associate 服务器在握手中途关闭 UnexpectedEof socks: server closed during the handshake
SocksUdpLink::associate 无法解析的回复 与解析器返回的相同 例如 unexpected server version
SocksUdpLink 控制 stream 读到 EOF,或控制 stream 读错误之后的调用 BrokenPipe socks: the control connection closed
Hysteria 2 拨号 服务器名称没有留下可用地址 AddrNotAvailable hysteria2: the server name resolved to no usable address
Hysteria 2 拨号 没有地址能完成连接 ConnectionRefused hysteria2: no address answered (<addr>: <error>; …)
Hysteria 2 拨号 CA 证书包不可读、含有不可用的证书,或不含证书 InvalidInput hysteria2: could not read the CA file: <error>、hysteria2: the CA file holds an unusable certificate: <error>、hysteria2: the CA file contains no certificates
Hysteria 2 拨号 槽正在退避 BrokenPipe hysteria2: connection is down, waiting before the next attempt
Hysteria 2 拨号 connect 任务没有给出结果就结束了 Other hysteria2: the connect task disappeared
Hysteria 2 UDP 打开 服务器不中继 UDP,或没有提供 QUIC 数据报 Unsupported hysteria2: the server does not relay UDP、hysteria2: the server did not offer QUIC datagrams
Hy2DatagramLink 接收 投递其回复的 channel 已关闭 BrokenPipe hysteria2: the connection behind this association is gone
WgDevice::start 没有隧道本地地址 InvalidInput wireguard: no tunnel-local addresses configured
WgDevice::start 端点解析不出任何结果 NotFound wireguard: endpoint domain did not resolve
WireGuard 拨号 槽正在退避 BrokenPipe wireguard: tunnel is down, waiting before the next attempt
WireGuard 拨号,解析 目的地解析不出任何结果 NotFound wireguard: destination did not resolve
WireGuard 拨号,解析 没有属于允许地址族的地址 AddrNotAvailable wireguard: no usable <strategy> destination address for <host>:<port>;隧道本地地址排除了候选地址时,追加 (local address supports IPv4 only) 或 (local address supports IPv6 only)
WireGuard TCP 拨号 隧道内没有地址能连上 TimedOut wireguard: tunnel TCP connect failed for all resolved addresses (<ip>: <error>; …)
Guarded 流的出站版本被排空策略关闭 ConnectionAborted the outbound this flow was opened on was drained
FanOutLink 子 link 的打开、发送或接收失败 无 记录日志;包被丢弃或子 link 被移除
级别 Target 文本
debug etemenanki_supervisor::topology::outbound::udp_fanout udp fan-out: opening <id> failed: <error>
debug etemenanki_supervisor::topology::outbound::udp_fanout udp fan-out: sending on <id> failed: <error>
debug etemenanki_supervisor::topology::outbound::udp_fanout udp fan-out: sub-link <id> ended: <error>
debug etemenanki_supervisor::topology::outbound::freedom freedom: dropping a datagram to an unresolvable <destination>
debug etemenanki_environment::dial::udp udp: no V4 socket: <error>、udp: no V6 socket: <error>
debug etemenanki_protocols::transports::keepalive could not enable TCP keepalive: <error>
info etemenanki_supervisor::topology::balancer balancer member <tag> is now up、balancer member <tag> is now down
warn etemenanki_protocols::hysteria::connector hysteria2: <error>; datagrams routed to this outbound are dropped(每个出站版本一次)
warn etemenanki_protocols::hysteria::slot hysteria2: connect failed: <error>、hysteria2: connection closed, reconnecting
warn etemenanki_protocols::wireguard::slot wireguard: tunnel start failed: <error>、wireguard: tunnel driver stopped, rebuilding

<id> 是子 link 的 OutboundId。Hysteria 2 和 WireGuard 的驱动还会记录它们自己的更多日志行,见 Hysteria 2:协议与客户端 和 WireGuard。

  • TCP 流的拨号 future 返回的 Err 在服务端核心中成为 Event::ConnectFailed。SOCKS 入站以失败回复应答它。
  • Guarded 或 Metered stream 在打开之后返回的 Err,在服务端运行时中是该 key 的出站错误(服务端运行时)。
  • fan-out 吞掉打开失败、发送失败和接收错误,所以 UDP 关联从不会从它那里看到出站错误。关联随其连接结束,如接收中所述。
  • 构建错误会拒绝这次应用,不做任何改变。

这一层的取消通过 drop 实现,探测和 Hysteria 2 的 connect 任务除外:

  • 在代理客户端的 StreamFuture 或 DatagramFuture 完成之前 drop 它,会 drop 传输连接和未完成的握手;对 SOCKS 关联还会 drop 控制 stream。对 freedom,它会放弃解析和连接尝试。
  • drop Hysteria 2 的拨号不会取消其连接尝试:start_connect 在独立的任务中运行 Hy2Conn::connect,即使所有等待方都被 drop,该任务也会完成,并把连接留在出站的槽中供下一个流使用(Hysteria 2:协议与客户端)。
  • 在 WgDevice::start 返回之后 drop WireGuard 拨号,会把已启动的隧道(包括驱动任务)留在出站的槽中(WireGuard)。
  • drop OutboundStream 会关闭 TCP stream,或 drop 客户端运行时及其上游 stream。
  • drop FanOutLink 会 drop 每个子 link 以及仍在打开中的子 link。这发生在连接的运行时结束时,每个子 link 的被跟踪流在它被 drop 时结束。
  • 探测任务在其负载均衡器被重建或移除时、以及关闭时,通过负载均衡器的 token 被取消。正在进行的探测被放弃,而不是等它完成。
  • supervisor 的出站代码只 spawn 探测任务。protocols crate spawn Hysteria 2 的 connect 任务,以及 Hysteria 2 连接和 WireGuard 隧道的驱动,见各自的协议页。
常量 位置 值 含义
MAX_SUBS supervisor/src/topology/outbound/udp_fanout.rs 64 每个 UDP 关联打开的子 link 数;最久未被发送的最先关闭
DEFAULT_PROBE_INTERVAL supervisor/src/topology/balancer.rs 30 秒 同一成员两次探测之间的休眠;probe_interval 未设置时 app 使用它
DEFAULT_PROBE_TIMEOUT supervisor/src/topology/balancer.rs 5 秒 单次探测的上限,包括解析和每次连接尝试;probe_timeout 未设置时 app 使用它
DEFAULT_CONNECT_TIMEOUT environment/src/dial/tcp.rs 10 秒 TcpDialer 的单次 TCP 连接尝试(见拨号器与 socket 策略)
TCP_KEEPALIVE_IDLE、TCP_KEEPALIVE_INTERVAL、TCP_KEEPALIVE_RETRIES protocols/src/transports/keepalive.rs 120 秒、30 秒、3 TransportConnector::dial 在每条代理上游 TCP 连接上设置的 keepalive 时间表
RECONNECT_BACKOFF_BASE、RECONNECT_BACKOFF_MAX protocols/src/hysteria/slot.rs 1 秒、30 秒 连续失败后 Hysteria 2 连接尝试之间的退避
REBUILD_BACKOFF_BASE、REBUILD_BACKOFF_MAX protocols/src/wireguard/slot.rs 1 秒、30 秒 连续失败后 WireGuard 隧道启动之间的退避
MIN_HEALTHY_LIFETIME protocols/src/hysteria/slot.rs、protocols/src/wireguard/slot.rs 10 秒 在此之前就失去的连接或隧道算作一次失败的尝试
TCP_CONNECT_ATTEMPT_TIMEOUT protocols/src/wireguard/slot.rs 10 秒 WireGuard 隧道内的单次 TCP 连接,每个解析出的地址一次
RECV_BUF protocols/src/socks/udp_link.rs 64 KiB 读取的最大 SOCKS 中继包;每个关联一个装箱缓冲区
sink protocols/src/socks/udp_link.rs 256 字节 读空 SOCKS 控制 stream 所用的缓冲区
scratch protocols/src/socks/udp_link.rs 2,048 字节 构建 SOCKS UDP 包所用的复用缓冲区的初始容量
HTTP_BUF、SOCKS_BUF、TROJAN_BUF、VLESS_BUF supervisor/src/topology/outbound/mod.rs 16,384 字节 客户端运行时缓冲区大小;每个流两个
VMESS_BUF supervisor/src/topology/outbound/mod.rs 32,768 字节 VMess 客户端运行时缓冲区大小
SS_BUF supervisor/src/topology/outbound/mod.rs 20,480 字节 SsStream::BUF_SIZE
SS2022_BUF supervisor/src/topology/outbound/mod.rs 65,569 字节 Ss2022Stream::BUF_SIZE = MAX_RECORD_LEN

对每个 UDP 关联,MAX_SUBS 限制打开的子 link 数。Trojan、VLESS 或 VMess 子 link 是一条经由传输层的独立上游连接,带两个客户端缓冲区;SOCKS 子 link 是一条独立的控制连接,加上一个中继 socket 及其 64 KiB 接收缓冲区。整个 workspace 的限制表见限制、超时与内存。

supervisor 的单元测试是库中以 #[path] 引入的模块,位于 supervisor/tests/unit/ 下;它的集成测试(socket_policy、hot_swap、tracking)在进程内运行 supervisor,对接回环地址上的 echo 服务器。app 的端到端测试是一个集成测试 crate app/tests/integration.rs,它启动真实的二进制。SocksUdpLink 的测试在 protocols 的 pipeline 测试 crate 中。

终端窗口
cargo test -p etemenanki-supervisor --lib topology::balancer
cargo test -p etemenanki-supervisor --lib topology::plane
cargo test -p etemenanki-supervisor --lib topology::spec_plan::plan
cargo test -p etemenanki-supervisor --lib build::validate
cargo test -p etemenanki-supervisor --test socket_policy
cargo test -p etemenanki-supervisor --test hot_swap
cargo test -p etemenanki-supervisor --test tracking a_udp_association
cargo test -p etemenanki-app --test integration e2e_udp_route
cargo test -p etemenanki-app --test integration e2e_balancer
cargo test -p etemenanki-protocols --test pipeline pipeline::socks::
测试 文件 验证的行为
failover_takes_the_first_healthy_member_in_order supervisor/tests/unit/balancer.rs a 不可用时选 b;a 恢复后再次选 a。
round_robin_cycles_only_through_healthy_members supervisor/tests/unit/balancer.rs b 不可用时,四次选择得到 a, c, a, c。
every_member_down_still_selects_rather_than_dropping supervisor/tests/unit/balancer.rs 全部不可用:仍返回第一个成员。
round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down supervisor/tests/unit/balancer.rs 首选 b 时四次都得到 b;b 不可用后,每次选择都是 a 或 c。
failover_ignores_the_preference supervisor/tests/unit/balancer.rs failover 下首选 b 仍得到 a。
a_preference_for_another_version_is_no_preference supervisor/tests/unit/balancer.rs 首选 b 的版本 7 时,结果按 a, b, a, b 交替。
a_balancer_needs_at_least_one_member supervisor/tests/unit/balancer.rs Balancer::new 拒绝空列表。
strategy_names_are_validated supervisor/tests/unit/balancer.rs failover 和 round_robin 能解析;roundrobin 不能。
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers supervisor/tests/unit/plane.rs 关闭负载均衡器的目标不影响经负载均衡的 stream;关闭成员的目标会中止它。
closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive supervisor/tests/unit/plane.rs 被排空的数据报 link 发送失败,已在等待的接收也失败;被排空的 fan-out 子 link 就是这样结束的。
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is supervisor/tests/unit/plan.rs 修改非成员时负载均衡器保持 v1;修改成员则构建 v2 并重新发布 plane。
removing_a_balancer_alone_republishes_the_plane supervisor/tests/unit/plan.rs 只移除一个负载均衡器会发布 plane,不排空任何东西,并保留它的版本。
a_balancer_needs_a_member、a_balancer_member_needs_an_upstream_a_tcp_probe_can_reach、a_balancer_member_must_be_a_known_outbound、a_balancer_member_may_not_be_another_balancer、a_balancer_may_not_share_an_outbounds_tag supervisor/tests/unit/validate.rs 哪些出站可以做负载均衡中的负载均衡器规则;探测规则针对 freedom、blackhole、Hysteria 2 和 WireGuard 检查。
direct_flows_are_dialed_on_hooked_sockets supervisor/tests/socket_policy.rs 直连 TCP 拨号是一个经过 hook 的 stream socket;直连 UDP 流至少使用一个经过 hook 的数据报 socket。
a_socks_upstream_is_dialed_on_hooked_sockets supervisor/tests/socket_policy.rs 经由 SOCKS 上游:CONNECT 拨号和关联的控制连接是两个经过 hook 的 stream socket,中继 socket 是一个经过 hook 的数据报 socket。
a_refusing_hook_fails_the_dial supervisor/tests/socket_policy.rs 返回 PermissionDenied 的 hook 使 SOCKS 请求失败。
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow supervisor/tests/hot_swap.rs 默认路由变为 block 后,活跃关联的下一个包被丢弃,而已打开的 TCP 流继续 echo。
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close supervisor/tests/hot_swap.rs freedom 出站 direct 的版本 1 以 Keep 排空,它的流继续运行;版本 2 以 Close 排空,它的流被关闭;新流使用最新版本。
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply supervisor/tests/hot_swap.rs 被复用的 Hysteria 2 出站不增加服务器会话;变化了的出站会拨号建立第二条连接,第一个流继续工作。
a_udp_association_counts_each_sub_link_and_charges_its_session supervisor/tests/tracking.rs 发往两个出站的三个包是两个流,字节数精确,负载计入 SOCKS 会话的 wire。
one_association_routes_each_peer_separately app/tests/integration/e2e_udp_route.rs 一个 SOCKS 关联:允许的对端应答,被端口规则屏蔽的对端不应答,之后允许的对端仍然应答。
replies_from_several_peers_merge_back_correctly app/tests/integration/e2e_udp_route.rs 一个关联上的两个对端各自得到自己的回复,并归属到正确的来源。
network_separates_tcp_from_udp app/tests/integration/e2e_route_context.rs network = "udp" 规则屏蔽 UDP 包,不影响 TCP。
traffic_moves_off_a_member_that_stops_answering app/tests/integration/e2e_balancer.rs probe_interval = 1 时,第一个成员停止接受 TCP 后,流离开它。
a_member_with_no_upstream_is_refused app/tests/integration/e2e_balancer.rs --test 拒绝基于 freedom 出站的负载均衡器。
server_runtime_relays_through_a_client_runtime_in_one_task、upstream_dial_failure_is_connect_failed_not_connected、refused_upstream_handshake_is_connect_failed_not_connected concepts/tests/client.rs 客户端运行时在同一个任务中作为服务端运行时的出站;拨号失败或握手被拒绝都是 ConnectFailed。
new_server_vs_new_client_udp protocols/tests/pipeline/socks.rs SocksUdpLink 经由 SOCKS 服务端核心往返数据报。
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 仅 Linux:绑定在 [::]:0 上的 link 能从 IPv4 映射地址收到 IPv4 中继的回复。
app_client_ws_xray_server_plain 及其他 app_client_* 测试 app/tests/integration/e2e_xray.rs VLESS 客户端对接真实的 Xray 服务器,覆盖 WebSocket、gRPC 和 TLS;没有 go 时跳过。
app_client_vmess_ws_xray_server_early_data_plain app/tests/integration/e2e_xray_vmess.rs VMess 客户端经由带 early data 的 WebSocket 对接 Xray;没有 go 时跳过。
app_socks_to_hysteria2_outbound、app_socks_to_hysteria2_with_salamander、app_socks_udp_to_hysteria2_outbound app/tests/integration/e2e_hysteria.rs 经由 Hysteria 2 出站的 TCP、Salamander 和 UDP,对接用 vendored 源码树构建的上游服务器;没有 go 时跳过。
app_socks_to_wireguard_outbound_tcp app/tests/integration/e2e_wg.rs 经由 WireGuard 出站到测试内对端的 TCP。

fan-out 的几项性质没有专门的测试:MAX_SUBS 淘汰、打开或发送失败时的丢包,以及 FanOutLink 内部逐包的负载均衡器首选(首选本身在 Balancer 上有单元测试)。socket 策略没有在 Hysteria 2、WireGuard 或探测的 socket 上测试。如果修改 FanOutLink,缺少的是一个基于手工构建的 Plane、带有 blackhole 和脚本化目标的单元测试;见测试。