跳转到内容

Connector 与 UDP fan-out

源码文件:24 个 · 核对版本 katana v3.0.1 · Etemenanki 596916d
  • katana/src/connector.rs
  • katana/src/router.rs
  • katana/src/rule.rs
  • katana/src/meter.rs
  • katana/src/serve.rs
  • katana/src/manager/transport.rs
  • katana/src/outbound/mod.rs
  • katana/src/outbound/proxy.rs
  • katana/src/outbound/freedom.rs
  • katana/tests/unit/connector.rs
  • katana/tests/unit/meter.rs
  • katana/tests/unit/rule.rs
  • katana/tests/unit/serve.rs
  • katana/tests/unit/e2e.rs
  • katana/tests/integration/sniff.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/helpers/address_family.rs
  • Etemenanki/protocols/src/vless/codec.rs
  • Etemenanki/protocols/src/vmess/codec.rs
  • Etemenanki/environment/src/routing.rs
  • Etemenanki/environment/src/dial/udp.rs

katana 节点上入站解码出的每条流,最终都汇入同一个调用:KatanaConnector::connect。一个 TCP 请求、每条 mux 子流、一个 UDP 关联以及每条 Hysteria 2 stream,都以 Flow<UserTag> 的形式到达这里。connector 按固定顺序同步地决定:用户能否打开这条流、由哪个出站承载、是否有审计规则禁止它,以及它的字节如何计费。它交给 per-connection 运行时的,要么是一条经过计量的字节流,要么是一个 FanOut,即逐个数据包路由 UDP 关联的数据报链路。

本页面向修改路由、审计、计费或出站代码的贡献者。本页先逐步讲解 src/connector.rs,然后深入介绍 FanOut,最后介绍直连出站的 UDP socket ResolvingUdp。准入本身、令牌桶以及驱动 connector 的服务端运行时各有单独的页面。

关注点 负责者 详见
该用户是否仍已注册,且仍对应这个 uid? Admission::admit,由 connector 调用 准入
告诉连接它承载的是谁的流 connector,通过 LeaseSlot 本页
选择出站 Router::route,使用由 route_target 构建的 RouteTarget 本页,以及路由模型
拒绝被审计的目标并记录命中 RuleManager::detect,通过 Dispatcher::forbidden 本页
拨号出站 Outbound::connect_stream 和 Outbound::connect_datagram 入站与出站
对字节计费,并在用户处于欠额时暂停传输 Gate 和 Metered 计量
对每个 UDP 数据包做路由、审计和计费 FanOut 本页
排队 effect、轮询链路、向协议核心报告失败 etemenanki-concepts 中的服务端运行时 服务端运行时

connector 不 spawn 任何 task,自己也不持有 socket。它返回的一切,包括 FanOut 及其内部的每条子链路,都由运行这条连接的运行时的那一个 task 轮询。丢弃运行时的出站槽位,就会把这些全部丢弃。

connector 为 etemenanki-protocols 中每个服务端协议核心产出的 Flow 类型实现了内核的 Connector trait。

concepts/src/link.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 trait DatagramLink: Unpin {
type Addr;
fn poll_send_to(
&mut self,
cx: &mut Context<'_>,
buf: &[u8],
to: &Self::Addr,
) -> Poll<io::Result<usize>>;
fn poll_recv_from(
&mut self,
cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<io::Result<Self::Addr>>;
}
pub enum Outbound<S, D> {
Stream(S),
Datagram(D),
}
protocols/src/flow.rs
pub struct Flow<T> {
pub destination: Destination,
pub user: NetworkUser<T>,
pub sniffed: Option<SniffedBehavior>,
pub source: Option<IpAddr>,
}

对于 UDP 流,destination 是第一个数据包的目标。后续数据包各自携带地址,作为 poll_send_to 的 to 参数到达数据报链路。只有当入站开启嗅探且目标是 IP 时,sniffed 才会被设置;见嗅探。

一个监听器的所有连接共享同一个 Dispatcher。src/manager/transport.rs 中的 TransportManager::start 在创建监听器时一并构建它,其中的 router 从不替换:修改路由会替换节点的传输层,连同其上的连接一起(见 tests/unit/e2e.rs 中的 route_change_drops_connections)。

src/connector.rs
pub struct Dispatcher {
pub router: Arc<Router<Outbound>>,
pub rules: Arc<RuleManager>,
pub node_tag: CompactString,
pub admission: Admission,
}
impl Dispatcher {
fn forbidden(&self, dest: &Destination, uid: i64) -> bool;
}
fn dest_string(d: &Destination) -> String;

forbidden 调用 RuleManager::detect(&self.node_tag, &dest_string(dest), Some(uid))。dest_string 给出审计规则所匹配的目标主机形式:域名或 IP 字面量,不带端口。

src/rule.rs
pub fn detect(&self, tag: &str, dest: &str, uid: Option<i64>) -> bool;

detect 以读模式获取规则表的 RwLock,按顺序尝试节点的规则,在第一个匹配的正则上,于 Mutex 保护下把 (uid, rule_id) 插入该节点命中记录的 HashSet。由于记录是一个集合,重复记录同一次命中不会新增任何内容,例如客户端不断向同一个被禁止的主机发送 UDP 数据包时。上报循环通过 RuleManager::drain 取走这个集合。

src/connector.rs
pub type LeaseSlot = watch::Sender<Option<CancellationToken>>;
#[derive(Clone)]
pub struct KatanaConnector {
disp: Arc<Dispatcher>,
source: Option<IpAddr>,
lease: Option<Arc<LeaseSlot>>,
}
impl KatanaConnector {
pub fn new(
disp: Arc<Dispatcher>,
source: Option<IpAddr>,
lease: Option<Arc<LeaseSlot>>,
) -> Self;
}
type ConnectFuture = Pin<
Box<dyn Future<Output = io::Result<link::Outbound<Metered<OutboundStream>, FanOut>>> + Send>,
>;
impl Connector<Flow<UserTag>> for KatanaConnector {
type Stream = Metered<OutboundStream>;
type Datagram = FanOut;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;
}

connector 克隆开销很小,每条连接一个。它在两个地方构建:

调用点 source lease 已退役用户的连接如何结束
src/serve.rs 中的 serve_stream(所有基于 TCP 的入站) 监听器看到的对端地址 Some:一个 watch channel 的发送端,其接收端由连接的 driver 监视 driver 等待已发布的租约,在租约被取消时结束整条连接
src/serve.rs 中的 run_hysteria Some(ip),客户端地址 None 监听器持有自己的运行时,所以每条流在其 Gate 拒绝传输字节时停止
src/outbound/mod.rs
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, dest: &Destination) -> StreamFuture;
pub fn connect_datagram(&self, dest: &Destination) -> DatagramFuture;
}
pub fn refused() -> io::Error;

refused() 即 io::Error::new(io::ErrorKind::PermissionDenied, "refused")。这条消息有意不给出原因,本文件中的每一处拒绝都使用它:从外部看,未知用户、被阻断的路由和审计命中完全一样。

src/outbound/proxy.rs 中的 OutboundDatagram 把出站打开的任意数据报链路装箱:

src/outbound/proxy.rs
pub struct OutboundDatagram(Box<dyn DatagramLink<Addr = Destination> + Send>);
impl OutboundDatagram {
pub fn new<D: DatagramLink<Addr = Destination> + Send + 'static>(link: D) -> Self;
}

两者都在计量中介绍。connector 这样使用它们:

src/meter.rs
impl Gate {
pub fn new(counter: Arc<UserCounter>, retired: CancellationToken) -> Self;
pub fn poll_open(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
pub fn sent(&self, n: usize);
pub fn received(&self, n: usize);
}
pub struct Metered<S> {
inner: S,
gate: Gate,
}
impl<S> Metered<S> {
pub fn new(inner: S, gate: Gate) -> Self;
}

当用户的租约未被取消、且其令牌桶没有欠额时,poll_open 就绪。租约一旦取消,它以 ConnectionAborted 失败,消息为 the user was retired。sent 和 received 累加用户的计数器,并向令牌桶扣费。

直到拨号之前的所有步骤都在 connect 内同步完成,早于返回的 future 第一次被轮询。只有拨号本身是异步的。

flowchart TB
  A["connect(flow)"] --> B{"admission.admit(tag)"}
  B -- None --> R1["Err(refused())"]
  B -- "Some(counter, lease)" --> C["发布租约(仅一次)"]
  C --> D["source = flow.source 或监听器看到的对端"]
  D --> E["Gate::new(counter, lease)"]
  E --> F{"destination.network"}
  F -- Udp --> G["Datagram(FanOut::new(...))"]
  F -- "其他" --> H["router.route(route_target(dest, sniffed, source))"]
  H --> I{"Block,或被审计禁止?"}
  I -- 是 --> R2["Err(refused())"]
  I -- 否 --> J["outbound.connect_stream(dest).await"]
  J --> K["Stream(Metered::new(stream, gate))"]
  1. 准入。 connect 克隆 flow.user.user_data 中的 Arc<UserTag>,并调用 self.disp.admission.admit(&tag)。只有当凭据仍注册在该 uid 名下时,它才返回用户的 Arc<UserCounter> 和租约。返回 None 时,connector 返回 Err(refused()) 且不记日志:已离开用户的客户端会不断重试,直到它察觉为止,每次尝试记一行会淹没日志。

  2. 只发布一次租约。 当 connector 持有 LeaseSlot 时,它调用 send_if_modified,传入的闭包只在槽位仍为 None 时存入租约,也只在此时报告发生了变化。因此,第一条通过准入的流告诉连接它属于哪个用户;之后的流既不会替换租约,也不会唤醒接收端。这一步发生在路由之前,所以即使连接的第一条流随后被拒绝,连接仍能获知自己的租约。

  3. 选择源地址。 flow.source.or(self.source):入站为这条流解码出的源地址优先于监听器看到的对端。

  4. 构建闸门。 Gate::new(counter, lease) 把用户的计数器和令牌桶与租约绑定。每条流有自己的 Gate;同一用户的多条流共享其背后的计数器和令牌桶。

  5. UDP:交出一个 FanOut。 当 flow.destination.network == DialNetwork::Udp 时,connector 立即返回 link::Outbound::Datagram(FanOut::new(disp, gate, tag.uid, source))。此时尚未路由、审计或拨号:关联没有单一的目标,所以这些都按数据包进行。

  6. TCP:路由。 其他网络类型走流路径。route_target(&flow.destination, flow.sniffed.as_ref(), source) 构建目标,Router::route 返回第一条匹配规则的出站,或默认出站。

  7. 拒绝。 如果出站是 Outbound::Block,或 forbidden(&flow.destination, tag.uid) 为 true,connector 返回 Err(refused())。这个判断是短路的:被 router 阻断的流永远不会拿去匹配审计规则,所以被阻断的目标不会记录审计命中。

  8. 拨号并计量。 outbound.connect_stream(&flow.destination) 开始拨号。返回的 future 等待拨号完成,用 ? 传播错误,并把流包装为 link::Outbound::Stream(Metered::new(stream, gate))。

src/router.rs
pub fn route_target<'a>(
dest: &'a Destination,
sniffed: Option<&'a SniffedBehavior>,
source: Option<IpAddr>,
) -> routing::RouteTarget<'a>;
environment/src/routing.rs
impl<Outbound> Router<Outbound> {
pub fn route(&self, target: RouteTarget<'_>) -> Arc<Outbound>;
}

route_target 先转换目标的 Remote 和端口,再补上手头已有的信息:

输入 Builder 调用 原因
sniffed: Some(s) with_sniffed_domain(s.domain.as_str()) 按 IP 寻址的流只在载荷中携带真实域名;有了它,域名和 geosite 规则才能匹配到这条流。
source: Some(ip) with_source(ip) 在监听器知道客户端地址时,供 RouteMatch::SourceCidr 使用。
dest.network 为 Tcp 或 Udp with_network(TargetNetwork::Tcp) 或 with_network(TargetNetwork::Udp) 供 RouteMatch::Network 使用;对应字段为 None 的匹配条件永远不会匹配。Unknown 和 Unix 不添加任何内容。

两个调用方传入的参数不同:

调用方 dest sniffed source
KatanaConnector::connect(流路径) flow.destination flow.sniffed.as_ref() 流的源地址,没有则用对端地址
FanOut::poll_send_to(每个数据包) 数据包的 to None 同一个源地址,在整个关联中固定

数据包从不嗅探,所以按 IP 寻址的 UDP 数据包只能匹配基于 IP 的规则(cidr、geoip)以及 port 规则。

一个 UDP 关联可以与多个对端通信,每个数据包各自指明对端。FanOut 就是 katana 为此交给运行时的数据报链路。它对每个数据包单独做路由、审计、闸门检查和计费,并为关联用过的每个出站保持一条打开的子链路。

src/connector.rs
pub const MAX_SUBS: usize = 64;
struct Sub {
key: usize,
link: OutboundDatagram,
}
pub struct FanOut {
disp: Arc<Dispatcher>,
gate: Gate,
uid: i64,
source: Option<IpAddr>,
subs: Vec<Sub>,
next: usize,
opening: Option<(usize, DatagramFuture)>,
recv_waker: Option<Waker>,
}
impl FanOut {
fn new(disp: Arc<Dispatcher>, gate: Gate, uid: i64, source: Option<IpAddr>) -> Self;
fn key_of(outbound: &Arc<Outbound>) -> usize;
fn poll_opening(&mut self, cx: &mut Context<'_>) -> Poll<()>;
}
impl DatagramLink for FanOut {
type Addr = Destination;
// poll_send_to, poll_recv_from
}
字段 含义
disp 当前 generation 的 dispatcher:router、审计规则和节点 tag。
gate 用户的限速与退役状态,每次发送和接收前都会检查。
uid 审计命中记录所对应的 uid。
source 每个数据包的路由目标中使用的客户端地址。
subs 已打开的子链路,最久未发送的在前,最近发送的在后。
next 下一次接收开始的下标,让每条子链路都能轮到。
opening 正在打开的那一条子链路,连同它的 key。
recv_waker 接收方的 waker,在没有子链路能交付数据时暂存于此。

子链路的 key 是 Arc::as_ptr(outbound) as usize。router 在整个 generation 期间持有每个出站的 Arc,而 FanOut 通过 disp 持有 router,所以在关联存续期间,这个指针可以稳定地标识“同一个出站”。

flowchart TB
  S["poll_send_to(buf, to)"] --> R["route(route_target(to, None, source))"]
  R --> X{"Block,或被禁止?"}
  X -- 是 --> D1["Ready(Ok(len)):丢弃,不计费"]
  X -- 否 --> G{"gate.poll_open"}
  G -- Pending --> P["Pending:运行时保留该数据包在队列中"]
  G -- Err --> E["Err(ConnectionAborted)"]
  G -- Ok --> L{"有该 key 的子链路?"}
  L -- 有 --> T["sub.poll_send_to,Ok(n) 时 gate.sent(n) 并把子链路移到末尾"]
  L -- 无 --> O{"opening"}
  O -- None --> N["opening = connect_datagram(to)"]
  N --> Q{"poll_opening,本 key"}
  O -- "本 key" --> Q
  O -- "其他 key" --> Q2{"poll_opening,其他 key"}
  Q -- Pending --> P
  Q -- 已打开 --> L
  Q -- 失败 --> D2["Ready(Ok(len)):丢弃,不计费"]
  Q2 -- Pending --> P
  Q2 -- "已打开或失败" --> L

用文字描述,poll_send_to 会:

  1. 用 route_target(to, None, self.source) 路由该数据包。
  2. 如果出站是 Block,或 forbidden(to, self.uid) 为 true,就返回 Poll::Ready(Ok(buf.len())) 丢弃它。UDP 没有办法只拒绝某一个对端,而关联还承载着其他对端,所以让发送失败是错误的。此时不检查闸门,也不计费。与流路径一样,被阻断的数据包不会拿去匹配审计规则。
  3. 等待 gate.poll_open(cx)。用户处于欠额时该调用返回 Pending,运行时把数据包保留在队首。每次轮询都从头开始,所以等待过的数据包在运行时重试时会被重新路由和审计。审计结果只在命中时记录,而记录是一个集合,所以客户端不断向同一个被禁止的主机发送时,在节点的上报循环用 RuleManager::drain 取走集合之前,只会新增一条 (uid, rule_id),而不是每个数据包一条。
  4. 查找与该出站 key 对应的子链路。如果有,就通过它发送。结果为 Ready(Ok(n)) 时调用 gate.sent(n),并把该子链路移到 subs 末尾。其他结果(Pending 或错误)原样返回。
  5. 没有匹配的子链路时,推进 opening:
    • None:启动 outbound.connect_datagram(to) 并继续循环。
    • 同一个 key:轮询这次打开。Pending 则保持 pending。成功时加入该子链路,循环后通过它发送。失败时以 Ok(buf.len()) 丢弃数据包,不计费,这样关联不会卡在一个无法到达的出站上。
    • 另一个 key:该数据包轮询那次打开直到结束(无论成功与否),然后循环去启动自己的打开。任何时候最多只有一次打开在进行。

poll_recv_from 也会推进进行中的打开,见下文“接收”。当接收侧完成了一次失败的打开时,等待中的数据包在下次轮询时会发现 opening 为空、也没有对应自己 key 的子链路,于是启动一次新的打开,而不是被丢弃。

poll_opening 轮询进行中的 DatagramFuture。成功时把新的 Sub 推到 subs 末尾。如果这使数量超过 MAX_SUBS,就移除 subs[0],即最久未发送的子链路,并用 saturating_sub(1) 把 next 减一,让轮转位置仍指向同一条子链路;如果 next 原本指向被淘汰的子链路,它保持为 0,指向新的队首。随后它唤醒 recv_waker,因为新子链路从未被轮询过回复,还没有在其 socket 上注册 waker。失败时它以 debug 级别记录 udp fan-out: opening an outbound failed: {e},并清空 opening。

stateDiagram-v2
  [*] --> Opening: 数据包被路由到一个没有子链路的出站
  Opening --> Open: connect_datagram 成功,推到末尾
  Opening --> [*]: connect_datagram 失败,数据包被丢弃
  Open --> Open: 发送成功,移到末尾
  Open --> [*]: 新子链路使数量达到 MAX_SUBS + 1 时被淘汰
  Open --> [*]: poll_recv_from 失败,被移除

失败的打开不会被记住。路由到同一出站的下一个数据包会再次尝试。淘汰一条子链路会丢弃它的链路:其 socket 或隧道随之关闭,仍发往它的回复会丢失。

各出站交给 connect_datagram 的结果:

Outbound 变体 子链路 数据包的寻址方式
Direct FreedomConnector::bind_udp,即 ResolvingUdp;不使用 dest 按数据包
Socks SocksUdpLink,通过一条到上游的新控制流发起 UDP ASSOCIATE;不使用 dest 按数据包
Wireguard WgConnector 为匿名流提供的数据报链路 按数据包
Vmess、Vless 朝 dest(即触发打开的那个数据包的目标)拨号的代理客户端运行时 发往打开时的目标
Http Err,Unsupported:http carries no datagrams 无
ShadowsocksLegacy、Shadowsocks2022 Err,Unsupported:shadowsocks carries no datagrams、shadowsocks-2022 carries no datagrams 无
Block 不会到达:FanOut 在打开任何东西之前就丢弃了数据包 无

对于 HTTP 和 Shadowsocks 出站,每个路由到它们的数据包都会耗费一次失败的打开尝试,然后被丢弃。

poll_recv_from:

  1. 等待 gate.poll_open(cx)。用户处于欠额时不读取任何子链路,回复留在子链路的 socket 或隧道中等待。
  2. 调用 poll_opening(cx),这样即使只有接收侧在被轮询,进行中的打开也能推进。
  3. 从 next 开始循环回绕,把每条子链路各轮询一次。第一条交付数据的子链路胜出:next 移到它之后的子链路,gate.received(n) 对它写入 buf 的字节计费,并返回它的发送方地址。
  4. 如果某条子链路失败,就以 debug 级别记录 udp fan-out: a sub-link ended: {e},移除该子链路;如果 next 此时越过末尾就重置为 0,再用 cx.waker().wake_by_ref() 唤醒自己,使剩余的子链路被重新轮询。错误不会返回:关联的生命周期长于其中任何一条子链路。
  5. 把 waker 存入 recv_waker 并返回 Pending。

从 next 而不是 0 开始,可以避免一条繁忙的子链路饿死其他子链路。

sequenceDiagram
  participant RT as 服务端运行时
  participant KC as KatanaConnector
  participant FO as FanOut
  participant OB as Outbound
  participant SL as 子链路
  RT->>KC: connect(network 为 Udp 的 flow)
  KC-->>RT: Datagram(FanOut)
  RT->>FO: poll_send_to(packet, to)
  FO->>FO: 路由、审计、gate.poll_open
  FO->>OB: connect_datagram(to)
  OB-->>FO: OutboundDatagram
  FO->>SL: poll_send_to(packet, to)
  SL-->>FO: Ok(n)
  FO->>FO: gate.sent(n)
  RT->>FO: poll_recv_from(buf)
  FO->>SL: poll_recv_from,从 next 开始轮转
  SL-->>FO: 回复,from
  FO->>FO: gate.received(n)
  FO-->>RT: Ok(from)
数据包的处理结果 是否计费 是否记录审计命中
被路由到 Block 否 否
匹配审计规则 否 是,(uid, rule_id) 记录一次
其子链路打开失败 否 否
被子链路接受 是,gate.sent(n) 否
被接受后在子链路内部丢弃,例如 ResolvingUdp 中没有可用地址的域名 是:子链路报告它已发送 否
从任一子链路读到的回复 是,gate.received(n) 否

Outbound::Direct 包装一个 FreedomConnector。它的 UDP 一侧是 ResolvingUdp,一个自行解析域名目标的双栈 socket,所以单条直连子链路就能服务关联指明的所有目标。

src/outbound/freedom.rs
const MAX_RESOLVED_NAMES: usize = 256;
impl FreedomConnector {
pub fn new(resolver: Resolver, strategy: AddressFamilyStrategy) -> Self;
pub async fn connect(&self, dest: &Destination) -> io::Result<TcpStream>;
pub fn bind_udp(&self) -> io::Result<ResolvingUdp>;
}
type ResolveFuture = Pin<Box<dyn Future<Output = Option<IpAddr>> + Send>>;
pub struct ResolvingUdp {
socket: DualStackUdp,
support: FamilySupport,
resolver: Resolver,
strategy: AddressFamilyStrategy,
resolved: HashMap<CompactString, Option<IpAddr>>,
order: VecDeque<CompactString>,
resolving: Option<(CompactString, ResolveFuture)>,
}
impl ResolvingUdp {
fn poll_target(&mut self, cx: &mut Context<'_>, to: &Destination) -> Poll<Option<SocketAddr>>;
}

bind_udp 调用拨号器的 bind_dual,为出站 address_family 允许的每个地址族各绑定一个 socket:

AddressFamilyStrategy 绑定的地址族
Ipv4Only V4
Ipv6Only V6
Auto、PreferIpv4、PreferIpv6 V4 和 V6

environment/src/dial/udp.rs 中的 bind_dual 在某个地址族绑定失败时以 debug 级别记录并继续绑定另一个;只有所有地址族都没绑定成功时才失败,错误为 AddrNotAvailable 和 udp: no usable local socket in any requested family。随后 bind_udp 在 FamilySupport 中记录实际绑定成功的 socket。只有这些地址族能发送,所以之后解析时会跳过没有 socket 的地址族的地址,改用解析结果中靠后的可用地址,而不是丢弃数据包。bind_dual 返回错误时打开失败,FanOut 丢弃该数据包。

poll_target 返回 Ready(Some(addr))、表示“丢弃”的 Ready(None),或 Pending:

  • IP 字面量在 strategy.allows(ip) 且 support.supports(ip) 时原样使用,否则为 None。
  • 已在 resolved 中的域名返回缓存的结果,可能是 None。
  • 新的域名在没有进行中的查询时启动一次查询:resolve_candidates("freedom", &dest, strategy, support, &resolver),取第一个候选地址。候选地址已按策略和已绑定的地址族过滤,并按 prefer_* 策略排序。查询出错则结果为 None。
  • 另一个域名的查询正在进行时,数据包需要等待:poll_target 返回 Pending,且不轮询那次查询。每条子链路同一时间只运行一次查询。只有当属于该域名的数据包被轮询时,查询才会推进,运行时会重试该 key 的 effect 队列队首的数据包来做到这一点。

查询完成后,结果插入 resolved,域名推入 order。如果 order 已存有 MAX_RESOLVED_NAMES 个域名,就先移除最早插入的那个。淘汰是先进先出的:缓存命中不会刷新域名的位置。

这个缓存本身没有过期机制。一个域名会一直保持它第一次得到的结果(包括否定结果),直到后续 256 个域名把它挤出,或子链路关闭。底层共享的 Resolver 有自己受 TTL 约束的缓存;见 DNS。

poll_send_to 通过双栈 socket 发往选定的地址。当 poll_target 返回 None 时,它以 debug 级别记录 freedom: dropping a datagram with no usable address 并返回 Ok(buf.len())。poll_recv_from 从任一 socket 读取,并把发送方报告为 Destination::udp(addr)。

不变量 机制 由谁固定
注册表中不存在的凭据,或注册在其他 uid 名下的凭据,什么也打不开 Admission::admit 首先查找 (key, uid);None 返回 refused() tests/unit/connector.rs 中的 a_user_the_registry_does_not_know_is_refused、a_credential_rebound_to_another_uid_is_refused
连接从第一条通过准入的流获知其用户,用户离开时该租约被取消 LeaseSlot::send_if_modified 只写入空槽位;Admission::commit 取消租约 tests/unit/connector.rs 中的 the_lease_reaches_the_connection_and_goes_with_the_user
已退役用户的连接会结束 src/serve.rs 中的 driver 等待已发布的租约 tests/unit/serve.rs 中的 a_retired_users_connection_ends
已退役用户不再传输任何字节,即使读操作已挂起 Gate::poll_open 首先轮询租约的 cancelled_owned future tests/unit/meter.rs 中的 a_retired_users_stream_refuses_to_move、retiring_the_user_wakes_a_parked_read
流在两个方向上都计入其用户 Metered 在每次传输完成后调用 Gate::sent 和 Gate::received tests/unit/connector.rs 中的 an_admitted_stream_is_billed_to_its_user;tests/unit/meter.rs 中的 each_direction_is_billed_to_the_user
限速与单次传输大小无关 每次传输后按完整字节数向令牌桶扣费,处于欠额时阻止后续传输 tests/unit/meter.rs 中的 the_limit_holds_however_the_writes_are_sized、the_limit_is_shared_by_both_directions
被 router 阻断的 TCP 流以 PermissionDenied 拒绝 拨号前判断 matches!(*outbound, Outbound::Block) tests/unit/connector.rs 中的 a_blocked_destination_is_refused
发往被审计目标的 TCP 流被拒绝,命中记录在该用户名下 Dispatcher::forbidden → RuleManager::detect tests/unit/connector.rs 中的 a_forbidden_destination_is_refused_and_recorded;tests/unit/rule.rs 中的 detect_records_and_drains
每个 UDP 数据包单独路由;被阻断的数据包被丢弃且不计费;已送达的数据包双向计费 FanOut::poll_send_to 中按数据包 route;只有子链路接受后才调用 gate.sent tests/unit/connector.rs 中的 udp_is_billed_after_routing_and_blocked_packets_are_free
发往被审计目标的 UDP 数据包被丢弃且不计费,命中被记录 在闸门之前按数据包调用 forbidden tests/unit/connector.rs 中的 a_forbidden_udp_destination_is_dropped_and_recorded
对按 IP 寻址的 TCP 流,嗅探出的域名能匹配到域名规则 route_target 添加 with_sniffed_domain tests/integration/sniff.rs 中的 a_sniffed_host_reaches_a_domain_rule、disable_sniffing_stops_the_domain_rule_matching
一个关联最多保持 MAX_SUBS 条子链路 推入后超过上限时,poll_opening 移除 subs[0] 没有专门的测试
一个关联同一时间最多只有一次打开在进行 opening 是单个 Option;其他 key 需等待它完成 没有专门的测试
一条直连子链路最多记住 MAX_RESOLVED_NAMES 个域名,且同一时间只运行一次查询 ResolvingUdp::poll_target 中的 order 淘汰和单个 resolving 槽位 没有专门的测试
位置 发生什么 运行时如何处理
admit 返回 None Err(refused()),不记日志 向协议核心发送 Event::ConnectFailed
TCP 路由为 Block,或审计禁止 Err(refused()) 向协议核心发送 Event::ConnectFailed
connect_stream 失败 拨号的 io::Error,由 connect future 返回 向协议核心发送 Event::ConnectFailed
UDP 数据包被阻断或被禁止 Ready(Ok(buf.len())),丢弃 视为已发送
子链路打开失败 debug 日志,Ready(Ok(buf.len())),丢弃 视为已发送
子链路发送失败 错误从 poll_send_to 返回 丢弃该数据包,保留关联,向协议核心发送 Event::SendFailed
子链路接收失败 debug 日志,移除该子链路,返回 Pending 无:关联在其他子链路上继续
用户已退役(发送时) Err(ConnectionAborted),the user was retired 丢弃该数据包,向协议核心发送 Event::SendFailed
用户已退役(接收时) Err(ConnectionAborted) 使整个出站 key 失败,向协议核心发送 Event::OutboundError
直连:数据包没有可用地址 debug 日志,Ready(Ok(buf.len())) 视为已发送

取消不需要额外代码。FanOut 持有它的子链路、进行中的 DatagramFuture 和 Gate;ResolvingUdp 持有它的 socket 和进行中的 ResolveFuture。当运行时丢弃出站槽位时(因为协议核心关闭了该 key、连接结束,或 generation 被拆除),丢弃这些对象会取消挂起的拨号和查询,并关闭所有 socket 和隧道。

名称 值 位置 作用
MAX_SUBS 64 src/connector.rs 一个 UDP 关联保持打开的子链路数。超过上限打开新子链路时,关闭最久未发送的那条。
进行中的打开 每个关联 1 个 FanOut::opening 发往另一个出站的数据包要等当前打开完成。
MAX_RESOLVED_NAMES 256 src/outbound/freedom.rs 一条直连 UDP 子链路记住的域名数。客户端使用更多域名时,被淘汰的域名需要重新查询。
进行中的查询 每条直连子链路 1 个 ResolvingUdp::resolving 发往另一个新域名的数据包要等当前查询完成。

本文件的单元测试位于 tests/unit/connector.rs,作为 tests 模块被 include 进 src/connector.rs。它们在一个包含 direct 和 block 出站的出站池之上构建真实的 Dispatcher,注册一个用户(uid 1),并可选地添加一条指向 block 的 CIDR 规则。测试使用真实的 loopback TCP 和 UDP echo 服务器,因此计费是按实际经过 socket 的字节核对的。

测试 固定的行为
a_user_the_registry_does_not_know_is_refused 未注册的凭据得到 PermissionDenied。
a_credential_rebound_to_another_uid_is_refused 现已注册到其他 uid 的凭据得到 PermissionDenied,因此新账户永远不会为旧账户的流付费。
an_admitted_stream_is_billed_to_its_user TCP 流以 Metered 流打开,完成 echo,上下行各计费 5 字节。
a_blocked_destination_is_refused 路由到 block 的目标得到 PermissionDenied。
a_forbidden_destination_is_refused_and_recorded 审计正则命中时拒绝该流,取出结果为 DetectResult { uid: 1, rule_id: 9 }。
the_lease_reaches_the_connection_and_goes_with_the_user 第一条流发布租约;提交空用户集合会取消租约;下一条流被拒绝。
udp_is_billed_after_routing_and_blocked_packets_are_free UDP 流以 FanOut 打开;送达的数据包及其 echo 各方向计费 4 字节;发往被阻断 CIDR 的数据包返回其长度且不增加计数。
a_forbidden_udp_destination_is_dropped_and_recorded 第一个目标被禁止时关联仍会打开;该数据包不计费;取出的命中为 DetectResult { uid: 1, rule_id: 4 }。

其他文件中的相关测试:tests/unit/meter.rs 覆盖 Gate 和 Metered,tests/unit/rule.rs 覆盖 RuleManager(detect_records_and_drains、update_skips_identical_ruleset、unknown_uid_records_negative_one),tests/unit/serve.rs 覆盖 a_retired_users_connection_ends。tests/integration/sniff.rs 中的嗅探测试让真实的 Xray 客户端连接 katana 节点,没有可用的 Xray 二进制时提前返回。

MAX_SUBS 淘汰、单个进行中的打开以及 ResolvingUdp 都没有专门的测试。修改其中任何一项时都应补充测试,使用一个 connect_datagram 可以保持 pending 或被设为失败的桩出站。