plane:为每个流选路
源码文件:39 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/router.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/build/route.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/topology/spec_plan/route.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/policy.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/topology/outbound/mod.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/topology/outbound/dns.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/environment/src/routing.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/protocols/src/socks/server.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/concepts/src/link.rsEtemenanki/concepts/src/runtime.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/router.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/app/tests/integration/e2e_route_context.rsEtemenanki/app/tests/integration/e2e_sniff.rsEtemenanki/app/tests/support/mod.rskatana/src/rule.rskatana/src/manager/node.rs
supervisor(监管器)承载的每个流都在 plane(数据平面)上选路。plane 是一份不可变的快照,包含编译好的路由表、路由规则所指的出站和负载均衡器目标,以及端口 53 的流要交给的 DNS 服务。这份快照放在唯一一个 ArcSwap cell 之后。一次改变路由、出站或 DNS 的应用(apply)会构建一份新快照,并用一次指针写入把它发布出去;每个 connector(连接器)对每个新流都会重新读取这个 cell。这就是 read-copy-update(RCU),也正因为如此,路由变更从不需要停止监听器,也不需要触碰任何已经打开的连接。
本页面向修改 supervisor/src/connector.rs、supervisor/src/topology/{plane,router,flow}.rs 或 supervisor/src/build/route.rs 的贡献者。内容包括:plane 及其 cell;路由表如何编译成槽位;流携带的类型(Flow、Principal、FlowContext);逐步讲解 AppConnector::connect;目标,以及让排空能够结束一个流的 Guarded 包装;epoch(plane 的发布序号)。一次应用如何决定是否发布,以及各种排空策略,见规划与应用变更;匹配条件本身见路由模型;目标所包装的出站、UDP 扇出和负载均衡器见出站、UDP 扇出与负载均衡器。
plane 和 connector 负责:
- 持有路由、目标和 DNS 服务的一份一致视图,并整体替换(
PlaneCell、Plane); - 在做任何其他事之前,先在流所属的会话上对它做准入(
Session::admit); - 每个 TCP 流只选路一次,使用调用
connect时的当前 plane;UDP 关联则逐包选路,使用每次调用FanOutLink::poll_send_to时的当前 plane; - 当 supervisor 有自己的 DNS 服务时,把端口 53 拦截给它,TCP 和 UDP 都一样;
- 每个流把负载均衡器解析为一个成员,使这个流记在实际承载它的出站名下,并随该出站一起排空;
- 在目标上打开流,并把它包装起来,使排空能够结束它(
Target::connect_stream、Guarded); - 把打开的流注册到 tracker(流跟踪器),附带它的规则、出站版本和 plane epoch(
Tracker::meter_stream)。
它们不负责:
- 计算匹配条件。
CompiledRoutes::decide把一个RouteTarget交给etemenanki_environment::routing::Router,由它按首个匹配遍历规则(路由模型)。 - 自己拨号。拨号由目标背后的
Outbound完成(出站、UDP 扇出与负载均衡器)。 - 决定何时发布 plane、哪些版本要排空。这由规划器和 actor 的提交(commit)决定(规划与应用变更)。
- 嗅探。协议核心、mux 解复用器和 SOCKS 驱动器在打开流之前填好
Flow::sniffed(嗅探)。 - 记录日志。
connect和 plane 不写任何日志行,TCP 流的失败以io::Error的形式返回给调用方。UDP 扇出以 debug 级别记录其子链路的失败(见日志行和出站、UDP 扇出与负载均衡器)。
katana 依赖「流记在解析后的出站名下」这条规则。它的目标审计把每条面板规则变成一条指向 tag 为 audit#<id> 的 blackhole 出站的路由,订阅 tracker,并对收到的每个 outbound().tag 能解析为 audit#<id> 的 FlowEvent::Opened 记录一次命中(src/rule.rs → audit_hit,src/manager/node.rs → spawn_audit)。审计规则只匹配域名,包括请求的域名和嗅探得到的域名。
谁读取 plane
Section titled “谁读取 plane”| 读取方 | 读取 cell | 选路方式 | 用途 |
|---|---|---|---|
AppConnector::connect(supervisor/src/connector.rs) |
每次调用一次,针对每个不是 UDP 的流(实际上就是 TCP) | Plane::route |
通过 connector 打开的每个 TCP 流:一条连接的唯一流、每个 mux 子流、每个 Hysteria 2 代理 stream、每个 SOCKS CONNECT。UDP 流只经过 connect 一次,变成一个扇出,不读取 cell。 |
FanOutLink::poll_send_to(supervisor/src/topology/outbound/udp_fanout.rs) |
每次调用 | Plane::route |
每个 UDP 关联的数据包 |
RoutedDialer::dial(supervisor/src/topology/outbound/dns.rs) |
每条到上游 DNS 服务器的连接一次 | Plane::route_upstream |
分流(split)解析器经代理发出的查询 |
Supervisor::epoch(supervisor/src/supervisor.rs) |
每次调用一次 | 无 | 报告当前 plane 的 epoch |
运行时在应用 Effect::Open 时同步调用 AppConnector::connect(concepts/src/runtime.rs;服务端运行时)。mux 载体的协议核心为每个子流推入一个 Open,所以每个子流都会经过 connect,TCP 子流会再次读取 plane。SOCKS 驱动器(protocols/src/socks/server.rs)不运行运行时:它自己调用 connect 并 await 返回的 future,每个 CONNECT 一次;每个 UDP ASSOCIATE 也是一次,发生在它转发该关联的第一个数据报时(SOCKS)。
PlaneCell 与 Plane
Section titled “PlaneCell 与 Plane”pub const DNS_PORT: u16 = 53;
pub type PlaneCell = Arc<ArcSwap<Plane>>;
pub struct Plane { routes: Arc<CompiledRoutes>, /// The target of each route slot, in `CompiledRoutes::targets` order. targets: Box<[Arc<Target>]>, /// The service applications' DNS is answered by, when the supervisor /// answers it itself. dns: Option<Arc<Target>>, epoch: u64,}
impl Plane { pub(crate) fn new( routes: Arc<CompiledRoutes>, targets: Box<[Arc<Target>]>, dns: Option<Arc<Target>>, epoch: u64, ) -> Self; pub(crate) fn empty() -> Self; pub fn epoch(&self) -> u64; pub fn route(&self, flow: &Flow, ctx: &FlowContext) -> Routed<'_>; pub fn route_upstream(&self, flow: &Flow, ctx: &FlowContext) -> Routed<'_>;}PlaneCell 是一个 Arc<arc_swap::ArcSwap<Plane>>(arc-swap 1.9)。每个 supervisor 只有一个 cell:Actor::new 用 ArcSwap::from_pointee(Plane::empty()) 创建它,并放进 Shared,每个 accept 循环都会克隆 Shared;Supervisor 句柄保留一份克隆,供 epoch 使用。读取方调用 load(),它不加锁,直接返回一个 guard;actor 调用 store() 来发布。Plane 构建之后永不修改。
targets按槽位索引。Plane::new用debug_assert_eq!检查它对routes的每个槽位恰好有一项。- 只有当 DNS 的 spec(期望状态)是
DnsSpec::Split时,dns才是Some;这是 supervisor 自己应答应用程序查询的唯一模式(名称解析与 DNS 服务)。 Plane::empty()是第一次提交之前 cell 中的值:一张只有一个槽位的表,槽位对应空 tag,由 id 为@v0的Outbound::Blackhole填充;没有 DNS 目标;epoch 为 0。the_empty_plane_drops_everything验证它会丢弃每一个流,包括端口 53 的流。Supervisor::start在返回之前应用第一个 spec,而那次提交在启动任何监听器之前就写入了第一个真正的 plane,所以 supervisor 自己的连接永远不会读到空 plane。
route、route_upstream 与 Routed
Section titled “route、route_upstream 与 Routed”#[derive(Clone, Copy)]pub struct Routed<'a> { pub target: &'a Arc<Target>, /// The rule that matched, if one did. pub rule: Option<RuleId>,}#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]pub struct RuleId(u32);
impl RuleId { pub const fn new(raw: u32) -> Self; pub const fn get(self) -> u32;}应用程序的流经过的是 route:
- 如果 plane 有 DNS 目标,并且
flow.destination.port == DNS_PORT,它返回该目标,并带上rule: None。这个检查只看端口,不看地址,也不看网络类型。代码注释给出了原因:UDP 应答被截断时,应用程序会改用 TCP 重试,而这次重试必须到达同一个服务。 - 否则调用
route_upstream。
route_upstream 向路由表要一个 Decision,返回 &self.targets[decision.slot] 以及 decision.rule。它从不拦截。DNS 服务自己的上游连接就用它选路,因为拦截这些连接会把服务的查询送回服务自己(port_53_is_intercepted_except_for_the_services_own_queries)。
RuleId 是匹配规则在为该流选路的那个 plane 的 RouteSpec::rules 中的下标。一次改变规则的应用之后,同一个数字可能指向另一条规则;流的 plane epoch 说明它索引的是哪个 plane 的规则。
Routed 从 plane 借用它的目标,所以它的生命周期不会超过 load() 返回的 guard。每个调用方都在 guard 被 drop 之前克隆 Arc<Target>:connect 通过 Target::resolve,扇出在解析之前用 routed.target.clone(),RoutedDialer 用 .target.clone()。
Target 与 TargetKind
Section titled “Target 与 TargetKind”pub struct Target { id: OutboundId, kind: TargetKind, /// Cancelled to close the flows opened on this version, when its drain /// policy says so. closed: CancellationToken,}
pub enum TargetKind { Outbound(Outbound), /// Chooses among outbound targets by health, per flow. Balancer(Arc<Balancer>),}
impl Target { pub fn outbound(id: OutboundId, outbound: Outbound) -> Self; pub fn balancer(id: OutboundId, balancer: Arc<Balancer>) -> Self; pub fn id(&self) -> &OutboundId; pub fn kind(&self) -> &TargetKind; pub fn resolve(self: &Arc<Self>, prefer: Option<&OutboundId>) -> Arc<Target>; pub fn connect_stream(self: &Arc<Self>, flow: Flow) -> impl Future<Output = io::Result<Guarded<OutboundStream>>> + Send + 'static; pub fn connect_datagram(self: &Arc<Self>, flow: Flow) -> impl Future<Output = io::Result<Guarded<OutboundDatagram>>> + Send + 'static; pub(crate) fn close_flows(&self);}
fn nested_balancer(id: &OutboundId) -> io::Error;一个目标是一个 tag 的一个版本:按其 spec 构建的一个出站,或者建立在若干出站目标之上的一个负载均衡器。每个目标都有自己的 CancellationToken,随目标一起创建,所以在 direct 版本 2 上打开的流可以被关闭,而不影响版本 3 上的流。TargetKind 带有 #[allow(clippy::large_enum_variant)]:目标在整个生命周期内都位于 Arc 之后,把较大的变体装箱只会多一层间接。
| 方法 | 作用 |
|---|---|
resolve(prefer) |
出站目标返回自身(一份 Arc 克隆)。负载均衡器返回 Balancer::select_preferring(prefer) 此刻选出的成员目标。prefer 是调用方在这条路由上已经在用的成员;在轮询策略下,只要它健康,负载均衡器就保留它,这样 UDP 关联就不会逐包在成员之间跳来跳去。connector 传入 None。 |
connect_stream(flow) |
在被调用时(而不是它的 future 第一次被 poll 时)以 None 解析,所以负载均衡器在调用时就选出成员。然后在该成员或出站上调用 Outbound::connect_stream(flow),克隆那个目标的 closed token,并返回一个把打开的 stream 包进 Guarded 的 future。 |
connect_datagram(flow) |
对 UDP 流做同样的事,底层是 Outbound::connect_datagram。扇出自己解析负载均衡器(它按成员的版本为子链路建 key),然后在成员上调用这个方法,此时 resolve 什么也不做。 |
close_flows() |
取消 closed。每个持有其克隆的 Guarded 在下一次 poll 时失败。只有 actor 的排空会调用它。 |
选择按流进行,而不是在选路时进行,所以被健康探测标记为不可用的成员从下一个流开始就会被跳过。成员如何探测、如何选择,以及所有成员都不可用时负载均衡器怎么做,见出站、UDP 扇出与负载均衡器。
如果解析得到的目标本身是负载均衡器,两个 connect 方法都会以 balancer member <tag>@v<n> is itself a balancer 失败(nested_balancer,一个 io::Error::other)。校验保证这种情况不会出现:负载均衡器的每个成员都必须是一个具有可探测上游的出站 tag。
Guarded
Section titled “Guarded”pub struct Guarded<T> { inner: T, closed: Pin<Box<WaitForCancellationFutureOwned>>,}
impl<T> Guarded<T> { fn new(inner: T, closed: CancellationToken) -> Self; fn poll_closed(&mut self, cx: &mut Context<'_>) -> Poll<io::Error>;}排空正是通过 Guarded 触及一个已经在中继的流。Guarded::new 持有内层 link,以及由目标 token 的克隆构建、装箱并 pin 住的 closed.cancelled_owned()。poll_closed 是各个受保护方法共用的唯一检查;token 被取消后,下面每个方法都会失败:
| Trait | 受保护的方法 |
|---|---|
AsyncRead |
poll_read |
AsyncWrite |
poll_write、poll_flush |
DatagramLink(Addr = Destination) |
poll_send_to、poll_recv_from |
token 被取消后,它们都返回 io::ErrorKind::ConnectionAborted,文本为 the outbound this flow was opened on was drained。挂起在一个没有动静的上游上的读或接收会被这次取消唤醒,并立即失败(closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive)。poll_shutdown 直接透传,不做这个检查。
对于运行时之下的 stream,运行时像对待任何其他出站错误一样处理这个错误:协议核心收到该 key 的 Event::OutboundError,关闭这个流。单流协议核心会结束整条连接;mux 载体只关闭那一个子流(服务端运行时)。SOCKS 中继从它所中继的 stream 上拿到这个错误。受保护的 UDP 子链路的错误留在它的 FanOutLink 内部,由后者丢弃这条子链路(出站、UDP 扇出与负载均衡器)。
CompiledRoutes 与 Decision
Section titled “CompiledRoutes 与 Decision”pub struct CompiledRoutes { table: routing::Router<Decision>, /// The tag each slot names, indexed by slot. targets: Vec<CompactString>,}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]pub struct Decision { /// The slot the flow is routed to. pub slot: usize, /// The rule that matched, or `None` for the default route. pub rule: Option<RuleId>,}
impl CompiledRoutes { pub fn new( rules: impl IntoIterator<Item = (Vec<routing::RouteMatch>, CompactString)>, default: CompactString, geo: GeoData, ) -> Self; pub fn single(tag: CompactString) -> Self; pub fn targets(&self) -> &[CompactString]; pub fn decide(&self, flow: &Flow, ctx: &FlowContext) -> Decision;}编译好的表指向的是槽位,而不是出站。规则或默认路由所指的每个不同 tag 各占一个槽位,plane 用同一次发布的目标填充这些槽位。把出站排除在表之外,才能在不重新编译表的情况下重建出站,并在出站变化时复用同一张编译好的表:plane 围绕同一个 Arc<CompiledRoutes> 重建,只换上新的槽位内容。
CompiledRoutes::new 按顺序遍历规则:
slot_of(tag)在一个HashMap<CompactString, usize>中查找 tag;第一次见到某个 tag 时,把它追加到targets并返回其下标。因此槽位按首次出现的顺序编号;如果没有任何规则提到默认路由的 tag,它最后才被加入。- 每条规则变成一个
RouteItem,其输出是它自己的Arc<Decision>,带有该规则的槽位和rule: Some(RuleId::new(index))。共用一个 tag 的规则共用槽位,但各自保留独立的 decision,所以一个流总能说明是哪条规则把它送过来的。下标用u32::try_from(at).expect("fewer than 2^32 rules")转换。 - 默认路由是带有其槽位和
rule: None的Arc<Decision>。 - 表本身是
routing::Router::new(RouteTable { routes, default }, geo)。
single(tag) 就是 new([], tag, GeoData::default()),Plane::empty 和单元测试用的就是这种表。decide 用 route_target 构建 RouteTarget,调用 Router::route,并把 Decision 从返回的 Arc 中复制出来。
对于规则 [port 80 → proxy, network udp → direct, port 443 → proxy]、默认 direct,every_rule_naming_a_tag_shares_that_tags_slot 验证了:
| 流 | targets() |
Decision |
|---|---|---|
| TCP 到端口 80 | ["proxy", "direct"] |
proxy 的槽位,规则 0 |
| TCP 到端口 443 | proxy 的槽位,规则 2 |
|
| UDP 到端口 443 | direct 的槽位,规则 1(首个匹配的规则胜出) |
|
| TCP 到端口 22 | direct 的槽位,None(默认路由) |
route_target
Section titled “route_target”pub fn route_target<'a>(flow: &'a Flow, ctx: &'a FlowContext) -> routing::RouteTarget<'a>;route_target 把 supervisor 的流适配为路由模型用来匹配的、与配置无关的目标。它借用所有字段,不做任何分配。
RouteTarget 字段 |
取自 | 读取它的匹配条件 |
|---|---|---|
remote |
flow.destination.remote,作为 Remote::Domain(&str) 或 Remote::IpAddr |
domain、GeoSite、Cidr、GeoIp |
port |
flow.destination.port |
PortRange |
sniffed_domain |
flow.sniffed 的 domain(如果已设置) |
所有域名匹配条件和 GeoSite,作为 remote 之外的另一个候选 |
inbound_tag |
总是 ctx.inbound_tag |
InboundTag |
source |
flow.source.or(ctx.source) |
SourceCidr |
network |
根据 flow.destination.network 取 TargetNetwork::Tcp 或 Udp;DialNetwork::Unknown 或 Unix 时不设置 |
Network |
这里有两个选择在代码中注明了原因:
- 流自己的源地址优先于监听器的。 代码注释给出了这个顺序的原因:用一个 socket 服务许多对端的监听器(例如 QUIC 监听器)只有在每个流上才知道对端是谁。在本版本的每条 accept 路径上,两者携带的都是同一个地址。
serve_connection把AppConnector::source()交给协议核心或 SOCKS 驱动器,由它们复制到自己打开的每个流上(Flow::new(dest, user, source));Hysteria 2 监听器为每条 QUIC 连接构建一个 connector,TUN 监听器为每个 TCP 流或 UDP 关联构建一个,每个 connector 携带的 IP 都与其协议核心放在流上的相同。只有 Unix socket 两者都不提供,而RoutedDialer的流没有源地址。supervisor/tests/unit/router.rs中的四个route_target测试覆盖了所有组合。 - 总是提供网络类型。 不提供的话,每条
RouteMatch::Network规则都将永远不起作用。
对应字段缺失的匹配条件不匹配。Unix socket 入站上的流没有源地址,所以没有任何 SourceCidr 规则会匹配它。
Flow、Principal 与 FlowContext
Section titled “Flow、Principal 与 FlowContext”pub type Flow = etemenanki_protocols::flow::Flow<Principal>;
pub struct Principal { user: Option<UserKey>, label: CompactString, revoked: OnceLock<UserRemovalPolicy>,}
impl Principal { pub fn user(user: UserKey, label: CompactString) -> Arc<Self>; pub fn anonymous() -> Arc<Self>; pub fn user_key(&self) -> Option<UserKey>; pub fn label(&self) -> &str; pub fn revoked(&self) -> Option<UserRemovalPolicy>; pub(crate) fn revoke(&self, policy: UserRemovalPolicy);}
pub fn anonymous_user() -> NetworkUser<Principal>;
#[derive(Clone, Debug)]pub struct FlowContext { pub inbound_tag: CompactString, pub source: Option<IpAddr>,}Flow 是 protocols crate 的流,以 Principal 作为其用户载荷:
pub struct Flow<T> { pub destination: Destination, pub user: NetworkUser<T>, /// What sniffing recovered from the flow's first bytes, when the inbound /// sniffs and the destination named only an IP. pub sniffed: Option<SniffedBehavior>, /// The client's address, when the inbound knows it. pub source: Option<IpAddr>,}
impl<T> Flow<T> { pub fn new(destination: Destination, user: NetworkUser<T>, source: Option<IpAddr>) -> Self; pub fn toward(&self, destination: Destination) -> Self;}NetworkUser<Principal> 携带客户端出示的授权信息和 user_data: Arc<Principal>。Flow::new 是协议核心、TUN 入站、SOCKS 驱动器和 RoutedDialer 使用的构造函数,初始为 sniffed: None。对 UDP 流而言,destination 是第一个数据包的地址;后续数据包各自携带自己的地址。Flow::toward(destination) 生成一个指向另一个目的地的流,用户和源地址不变,且 sniffed: None;扇出用它把每个数据包当作一个独立的流来选路,mux 解复用器用它打开每个子流。Clone 是手写实现的,这样克隆时从不要求 Principal: Clone。
Principal 表示流属于谁,即它的 principal(身份主体):入站准入这个流时认定的用户,或者无人。每个服务端协议的用户表在握手成功时都会交回一个 principal,所以会话从它打开的第一个流得知自己的用户。每个入站为每个用户持有自己的 principal,所以在一个入站上吊销某个用户(revoke,第一次调用生效,存放在一个 OnceLock 中)不会影响该用户在别处的会话。connector 两次用到 principal:Session::admit 在第一个流上把会话绑定到它,并拒绝已被吊销的 principal;FlowMeta 保留它,让 tracker 能找到该用户的限速和标签。准入规则和吊销竞争见用户、principal 与会话。
Principal::anonymous() 没有用户 key,标签为空。它用于处于开放模式或共享凭据模式的入站,以及 supervisor 自己打开的流。anonymous_user() 把它包进一个 NetworkUser,其授权信息是用户名和密码都为空的 UsernamePassword;RoutedDialer 和单元测试使用它。
FlowContext 描述流从哪里进来。它的两个字段之所以存在,只是因为规则可以匹配它们:携带这两个字段,RouteMatch::InboundTag 和 RouteMatch::SourceCidr 才不会成为永远无法触发的规则。source 是监听器在 accept 时捕获的对端地址(如果有的话)。
AppConnector 与 FlowScope
Section titled “AppConnector 与 FlowScope”#[derive(Clone)]pub struct AppConnector { plane: PlaneCell, ctx: FlowContext, tracker: Tracker, /// The session the flows belong to. A TUN inbound's flows have none. session: Option<Arc<Session>>, /// Whether datagram flows count toward the session's wire bytes. charge_datagrams: bool,}
impl AppConnector { pub(crate) fn new( plane: PlaneCell, ctx: FlowContext, tracker: Tracker, session: Option<Arc<Session>>, ) -> Self; pub(crate) fn charging_datagrams(self) -> Self; pub fn source(&self) -> Option<IpAddr>; pub fn session(&self) -> Option<&Arc<Session>>;}
// 这里的 `Outbound` 是 `etemenanki_concepts::link::Outbound`。type ConnectFuture = Pin< Box< dyn Future<Output = io::Result<Outbound<Metered<Guarded<OutboundStream>>, FanOutLink>>> + Send, >,>;
impl Connector<Flow> for AppConnector { type Stream = Metered<Guarded<OutboundStream>>; type Datagram = FanOutLink; type Future = ConnectFuture;
fn connect(&mut self, flow: Flow) -> ConnectFuture;}
#[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>>,}AppConnector 实现了 concepts/src/link.rs 中的 Connector<Flow> trait(fn connect(&mut self, target) -> Self::Future,resolve 为 link::Outbound::Stream 或 link::Outbound::Datagram)。只有本 crate 构造它。前端程序从不自己构建 connector,而是通过运行入站来获得 supervisor 的 connector。在 crate 内部,serve_connection 读取 source(),把它作为客户端地址交给协议核心或 SOCKS 驱动器;session() 由 serve_connection 的 SOCKS 分支读取(它把 stream 包进 WireCounted,把其字节计入会话的 wire,即会话在线上收发的字节),drive 也会读取它。
只有 SOCKS 会通过 charging_datagrams() 设置 charge_datagrams:SOCKS UDP 关联的数据报从不经过会话自己的 stream,所以改为把它们的载荷计入会话的 wire。其他所有入站的数据报都在会话的传输层内部传输,已经在那里计过了。wire 的说明见按用户的用量计费。
connector 从哪里来(accept 路径见监听器与服务循环):
| accept 路径 | 一个 connector 对应 | ctx.source |
session |
|---|---|---|---|
基于 TCP 的 stream 监听器(Connection::serve_socket) |
一个会话:accept 得到的 socket 的会话;在 gRPC 上则是每个 HTTP/2 stream 的会话,每个 stream 自成一个会话 | 对端的 IP | 该会话 |
| 基于 Unix socket 的 stream 监听器 | 一个 accept 得到的 socket | None |
该 socket 的会话 |
Hysteria 2(run_hysteria_inbound) |
一条 QUIC 连接 | 客户端的 IP | 该连接的会话。在监听器停止之后才完成握手的连接得到 None 和一个已经取消的 token,所以会立即被关闭,没有会话能逃过被移除入站的关闭。 |
TUN(run_tun_inbound) |
一个 TCP 流,或按源地址和端口划分的一个 UDP 关联 | 主机的 IP | None:设备不准入任何用户 |
stream 和 Hysteria 2 的 connector 通过 Shared(supervisor/src/serve.rs)构建。Shared::connector 用 Sessions::open 注册一个新会话并为它构建 connector;监听器的停止 token 被取消后,Sessions::open 返回 None,此时不构建 connector。Shared::connector_for 为已经注册的会话构建 connector,例如 stream accept 循环为每个 accept 得到的 socket 打开的会话。TUN 的 connector,以及停止之后才到达的连接的 Hysteria 2 connector,用 AppConnector::new 构建,不带会话。
plane 的结构
Section titled “plane 的结构”flowchart LR cell["PlaneCell (ArcSwap)"] --> plane["Plane,epoch N"] plane --> routes["Arc CompiledRoutes"] plane --> slots["targets,每个槽位一个"] plane --> dnst["dns: Option Target"] routes -- "Decision.slot" --> slots slots --> direct["Target direct@v2"] slots --> lb["Target lb@v1(负载均衡器)"] lb --> a["成员 a@v1"] lb --> b["成员 b@v3"] dnst --> svc["Target @vN,Outbound::Dns"]
负载均衡器目标通过它的各个 Member 持有成员的目标:这些 Arc<Target> 与 actor 为那些出站 tag 持有的是同一批;当某条规则也直接指向某个成员时,槽位持有的也是同一个。DNS 目标是内部的:它的 tag 为空,版本是它首次被发布时所在 plane 的 epoch。
为一个流选路
Section titled “为一个流选路”flowchart TB
start["Plane::route(flow, ctx)"] --> dns{"有 DNS 目标且端口为 53?"}
dns -- "是" --> svc["DNS 目标,rule None"]
dns -- "否" --> up["route_upstream"]
up --> rt["route_target(flow, ctx)"]
rt --> pick["Router::route,首个匹配"]
pick -- "规则 i 匹配" --> d1["Decision:slot,Some(i)"]
pick -- "没有规则匹配" --> d0["Decision:默认槽位,None"]
d1 --> slot["targets[slot]"]
d0 --> slot
slot --> res["Target::resolve"]
res --> member["出站目标(自身,或负载均衡器的一个成员)"]
一个 TCP 流经过 connect
Section titled “一个 TCP 流经过 connect”sequenceDiagram participant RT as 运行时或 SOCKS 驱动器 participant AC as AppConnector participant S as Session participant P as PlaneCell participant T as Target participant TR as Tracker RT->>AC: connect(flow) AC->>S: admit(principal) AC->>P: load() P-->>AC: 当前 Plane 的 guard AC->>AC: route、resolve(None)、FlowMeta::of AC->>T: connect_stream(flow) T-->>AC: future(出站拨号、closed token) AC-->>RT: 装箱的 future,guard 已释放 RT->>AC: poll 这个 future T-->>AC: Guarded stream AC->>TR: meter_stream(meta, stream) AC-->>RT: link::Outbound::Stream(Metered)
AppConnector::connect 逐步讲解
Section titled “AppConnector::connect 逐步讲解”- 准入。 当 connector 有会话时,最先运行
session.admit(&flow.user.user_data),TCP 和 UDP 都一样。拒绝以一个已经完成的 future 返回,所以调用方看到的是一次失败的 connect。规则见用户、principal 与会话;拒绝的文本见失败路径与取消。 - UDP 变成扇出。
destination.network为DialNetwork::Udp的流不在这里选路。connector 构建FanOutLink::new(plane, ctx, flow, FlowScope { tracker, session: session.id(), charge })(只有设置了charge_datagrams时,charge才是会话的 wire),并立即返回link::Outbound::Datagram(link)。此时还没有注册任何流;每条子链路在打开时各自注册。 - 读取 plane。 其他所有流,不论网络类型,都继续往下走:
let plane = self.plane.load();。代码中的注释说明了这一步实现的规则:为这个流只读取一次,所以路由变更会作用到下一个流,而这个流保留它打开时所在的目标。 - 选路。
plane.route(&flow, &self.ctx)返回目标和匹配的规则;端口 53 被拦截时返回 DNS 目标。 - 解析。
routed.target.resolve(None)。负载均衡器此时选出成员,使这个流记在实际承载它的出站名下。 - 描述。
FlowMeta::of(&flow, session id, &ctx.inbound_tag, ctx.source, target.id().clone(), routed.rule, plane.epoch())确定 tracker 将要报告的内容:会话、入站 tag、principal、源地址(flow.source.or(ctx.source))、请求的目的地、嗅探到的域名、出站版本、规则和 plane epoch。 - 打开。
target.connect_stream(flow)调用出站的connect_stream,并捕获目标closedtoken 的一份克隆。connect返回时 plane guard 被 drop;装箱的 future 携带出站自己的 future、token 克隆、一份Tracker克隆和FlowMeta。 - 计量。 调用方把 future poll 到完成时,
tracker.meter_stream(meta, stream)把打开的 stream 包装成Metered<Guarded<OutboundStream>>,注册这个流并发出FlowEvent::Opened。失败的拨号会在这一步之前通过?返回错误,所以打开失败的流永远不会被注册。
第 1 到 7 步在调用方的 task 中、在 connect 调用内同步执行:调用方是应用 Effect::Open 的运行时,或 SOCKS 驱动器。connect 不 spawn 任何东西,第 8 步在调用方 poll 返回的 future 时执行。
UDP:逐包选路
Section titled “UDP:逐包选路”FanOutLink::poll_send_to 每次被调用都会选路:self.flow.toward(to)、self.plane.load()、plane.route(...),然后用该关联在这条路由上已经在用的成员调用 resolve(prefer),最后使用以解析得到的 OutboundId 为 key 的子链路。子链路以 id 而不是 Arc<Target> 的地址为 key,原因是:出站的寿命超过单次应用,一个 tag 可能有一个正在排空的旧版本与新版本并存,而且被释放的内存可能被其后继复用(见 supervisor/src/entity/id.rs 和 udp_fanout.rs 的模块注释)。发往端口 53 的数据包与 TCP 流一样会被拦截,因为扇出调用的是同一个 route。
每条子链路都是一个独立的被跟踪流,其 FlowMeta 由 Flow::toward(to) 构建:
FlowEntry::destination()是打开这条子链路的那个数据包的目的地;FlowEntry::sniffed()总是None,因为toward会清除它;FlowEntry::rule()是为那个数据包选路的规则;走默认路由或由 supervisor 的 DNS 服务应答的查询为None;FlowEntry::plane_epoch()是为那个数据包选路的 plane 的 epoch。
a_udp_association_counts_each_sub_link_and_charges_its_session 验证:被路由到两个出站的两个目的地产生两个流,一个的 rule() 为 None,另一个为规则 0。子链路表、它的大小(MAX_SUBS = 64)以及发送和接收规则,见出站、UDP 扇出与负载均衡器。
DNS 拦截
Section titled “DNS 拦截”在 DnsSpec::Split 下,build_dns 返回一个 Service,准备(prepare)阶段把它包装成 Target::outbound(internal_id(epoch + 1), Outbound::Dns(service))。internal_id 构建一个 tag 为空的 OutboundId;校验拒绝任何 spec 中的空 tag,所以这个 id 永远不会与 spec 中的出站冲突。版本号是同一次提交即将写入的 epoch,所以每个重建的 DNS 目标都有一个此前任何 DNS 目标都没有用过的版本。
| 流 | 选路结果 | 打开方式 |
|---|---|---|
| TCP 到任意地址的端口 53 | DNS 目标,rule: None |
基于该服务上一个 DnsTcpStream 的 OutboundStream::Proxy,立即就绪,不拨号 |
| 发往任意地址端口 53 的 UDP 数据包 | DNS 目标,rule: None |
该服务上的一条 DnsUdpLink 子链路,立即就绪,不拨号 |
任意流、端口 53,DnsSpec::Single |
路由表,与任何其他端口一样 | 路由表指向的出站 |
解析器自己到上游服务器的连接(RoutedDialer) |
route_upstream:只走路由表 |
路由表为该服务器地址指定的出站 |
RoutedDialer 只为 through_proxy 列表非空的分流解析器构建,这种解析器用直连解析器做引导(supervisor/src/build/dns.rs)。它把每条上游连接当作一个指向服务器 IP 和端口的 TCP 流来选路,使用匿名用户和 FlowContext { inbound_tag: DNS_FLOW_TAG, source: None },其中 DNS_FLOW_TAG 为 "dns",所以规则可以通过入站 tag dns 匹配解析器的流量。它用 Target::connect_stream 打开连接,所以 stream 是 Guarded 的,关闭型排空能够触及它;它从不经过 tracker,所以这些流不会被注册,不能被强制关闭(kill),也不计入任何人的用量。它以 Weak<ArcSwap<Plane>> 的形式持有 cell,因为 plane 通过它的出站拥有解析器,而解析器又拥有这个拨号器。dial 同步升级这个引用;一旦没有任何东西再持有 plane cell(supervisor、它的 accept 循环、每个存活的 connector 和扇出 link),升级就会失败,返回的 future 在被 poll 时以 dns: the supervisor this resolver dials through is gone 失败。服务本身见名称解析与 DNS 服务和 DNS 解析器。
嗅探到的目的地与请求的目的地
Section titled “嗅探到的目的地与请求的目的地”每个做嗅探的打开方都在打开流之前设置 flow.sniffed:
| 打开方 | 嗅探到的名称来自 |
|---|---|
| 做嗅探的协议核心(HTTP、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022、Hysteria 2 stream 核心) | 它的嗅探前缀;在协议核心推入 Effect::Open 之前设置,位于其 open_sniffed 中,或在打开之前内联设置 |
PassthroughCore,服务 TUN 流(protocols/src/core/mod.rs) |
它的嗅探前缀,在 open 中 |
mux 解复用器(protocols/src/mux/demux.rs) |
子流 New 帧的载荷,经由 crate::sniff::sniff |
SOCKS 驱动器(protocols/src/socks/server.rs) |
它在批准请求之后、调用 connect 之前收集的前缀 |
plane 只把嗅探到的名称当作域名匹配条件的一个额外输入(RouteTarget::sniffed_domain)。connect 路径上没有任何东西改写 flow.destination:出站拨号的是客户端请求的地址,也没有任何出站读取 sniffed。FlowMeta 两者都保留,所以 FlowEntry::destination() 报告请求的地址,FlowEntry::sniffed() 报告嗅探到的名称。
app 的端到端测试让这一点可以被观察到。在 app/tests/integration/e2e_sniff.rs 中,通往 freedom 的唯一路由是一条针对 sniffed.example 的 domain_suffix 规则,默认路由是 blackhole,客户端请求 127.0.0.1,而 sniffed.example 无法解析。因此一次成功的回显同时证明了两件事:嗅探到的名称被匹配,并且流被拨号到了请求的 IP。嗅探器见嗅探。
pub(crate) fn compile_routes(spec: &RouteSpec) -> io::Result<CompiledRoutes>;pub struct RouteSpec { pub rules: Vec<RouteRuleSpec>, pub default: CompactString, pub geoip: Option<PathBuf>, pub geosite: Option<PathBuf>,}
pub struct RouteRuleSpec { pub matchers: Vec<RouteMatch>, pub outbound: CompactString,}规则的任一匹配条件匹配时,该规则就匹配;规则按顺序尝试;没有任何规则匹配的流走 default。outbound 和 default 都指向一个出站或负载均衡器:两者共用一个 tag 命名空间。geoip 和 geosite 指定 compile_routes 通过 build_geo_data 读取的文件。
compile_routes 在 Actor::prepare 中为计划里的 Step::Build(Resource::Route) 运行。check 是 etemenanki-app --test 背后的试运行,它从一个全新的 actor 出发做规划,所以总会编译路由。compile_routes 会:
-
按规则顺序收集所有规则中的每个
RouteMatch::GeoSite(code)和RouteMatch::GeoIp(code),包括重复项; -
调用
build_geo_data(spec.geoip, spec.geosite, &geosite_codes, &geoip_codes)(environment/src/routing.rs),它:- 先处理 geosite 再处理 geoip,所以 geosite 和 geoip 引用都失败的 spec 报告的是 geosite 的错误;
- 只有至少用到一个同类匹配条件时才读取对应的 geo 文件,所以没有匹配条件使用的已配置路径永远不会被打开(只配置了端口规则时,即使
geoip路径不存在,etemenanki-app --test也会打印Configuration OK.); - 每个不同的匹配字符串只加载一次,跳过已经在 map 中的,并以原始字符串作为结果的 key,例如
google@ads或!cn; - 查找代码时忽略 ASCII 大小写:geosite 取
@之前的部分,geoip 取开头!之后的部分;
解码方式以及
@attribute和!的作用见路由模型; -
把规则(作为
(matchers, outbound)对)、默认路由和 geo 数据交给CompiledRoutes::new。
任何错误都会让这次应用以 ApplyError::Build { resource: Resource::Route, .. } 失败,正在运行的 plane 保持原样。校验在此之前运行,所以指向既不是出站也不是负载均衡器的 tag 的规则或默认路由永远到不了这里。下面是 etemenanki-app --test 打印的文本,每条前面都带有 configuration invalid: 前缀:
| 原因 | 文本 |
|---|---|
| 规则或默认路由指向未知 tag(校验) | route references unknown outbound <tag> |
| 使用了 geosite 匹配条件但没有 geosite 文件 | building route failed: a geosite matcher is used but no geosite file is configured |
| 使用了 geoip 匹配条件但没有 geoip 文件 | building route failed: a geoip matcher is used but no geoip file is configured |
| 文件中没有该代码 | building route failed: geosite code not found: <code> 或 building route failed: geoip code not found: <code>。<code> 是去掉 @attribute 或开头 ! 之后的名称:nosuchcode@ads 和 !nosuchcode 都打印为 nosuchcode。 |
| 文件不是 geo 列表 | building route failed: geosite decode: <decoder error>(或 geoip decode: …) |
| 文件无法读取 | building route failed: <OS error>,例如 No such file or directory (os error 2) |
| 某个 geosite 代码的正则条目在跳过无效条目之后,合在一起仍然编译失败 | building route failed: geosite <code>: regex entries: <error> |
| geoip 条目的地址不是 4 或 16 字节 | building route failed: geoip cidr: ip must be 4 or 16 bytes, got <n> |
geoip 条目的前缀放不进 u8 |
building route failed: geoip cidr: prefix too large |
| geoip 条目的前缀对其地址来说太长 | building route failed: geoip cidr: <error> |
有两种 geosite 问题以 warn 级别(target 为 etemenanki_environment::routing)记录并被跳过,而不是被拒绝:未知类型的域名条目(skipping unknown geosite domain type <n>),以及无法编译的正则条目(skipping invalid geosite regex "<pattern>": <error>)。
无效的端口范围和域名正则更早就被拒绝:前端程序用 parse_port_match 和 parse_domain_regexes 把规则降为 spec 时就会拒绝它们(见etemenanki-app:从 TOML 到 spec)。
发布 plane
Section titled “发布 plane”当 DNS spec、某个出站(新增、修改、移除,或因 DNS 变更而重建)、某个负载均衡器(新增、修改、移除,或因其某个成员被重建而重建)或路由 spec 发生变化时,规划器加入 Step::PublishPlane。用户集和入站从不触发发布,通过 set_users、upsert_user 和 remove_user 做的用户编辑也从不经过规划器,不会碰 cell。每种情况以及验证它们的计划测试见规划与应用变更。
在 Actor::commit 中,先提交用户 key 和限速,并用准备阶段构建或复用的目标替换 self.targets。然后,仅当计划包含 Step::PublishPlane 时:
self.epoch += 1。- 填充槽位:按槽位顺序,对
routes.targets()中的每个 tag 取self.targets[tag].clone()。每个这样的 tag 都存在,因为校验已经拒绝了任何指向不存在 tag 的规则。 self.shared.plane.store(Arc::new(Plane::new(routes.clone(), slots, dns_target.clone(), self.epoch)))。从这次写入起,每次load()都会看到新的 plane。- 取消计划中每个不是
Reuse的负载均衡器的探测 token。对每个新的负载均衡器,创建根 token 的一个子 token,并在 supervisor 的TaskTracker上调用Balancer::spawn_probe:用新的dns.servers解析器解析,在 supervisor 的 socket 策略下用TcpDialer::new(socket)连接(出站、UDP 扇出与负载均衡器)。
每次提交,无论是否发布,actor 随后都会把 routes、dns 和 dns_target 保存为正在运行的值。没有 PublishPlane 时,cell、epoch 和探测都保持不变。
提交中的其余所有工作都在这次写入之后运行,包括启动新的监听器,计划中的 Drain 步骤排在最后。由此得出两个结果。一次应用所启动的监听器永远不会在上一个 plane 上服务连接。旧的出站版本只有在新流已经流向其后继版本之后才会被排空;a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish 断言 PublishPlane 在 Drain 步骤之前。提交以 debug 日志行 applied: <n> built, <n> reused, <n> swapped, <n> drained 结束。
被复用的出站或负载均衡器在旧 plane 和新 plane 中是同一个 Arc<Target>,token 也相同。未变化的有状态出站(例如带着 QUIC 连接的 Hysteria 2 客户端)正是这样跨越一次应用保留下来的(supervisor/tests/hot_swap.rs 中的 an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply)。
| 位置 | 含义 |
|---|---|
Plane::epoch() |
这个 plane 发布时的 epoch:Plane::empty() 为 0,第一个 spec 的 plane 为 1,之后每次发布加一。 |
Supervisor::epoch() |
self.plane.load().epoch():当前 plane 的 epoch。它直接读取 cell,不经过 actor。 |
FlowEntry::plane_epoch() |
为该流选路的 plane 的 epoch,来自 FlowMeta。对 UDP 子链路而言,是为打开它的那个数据包选路的 plane。 |
| DNS 目标的版本 | 该目标首次发布时所在 plane 的 epoch。 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow 断言修改默认路由会让 Supervisor::epoch() 恰好加一;a_refused_spec_changes_nothing 断言被拒绝的 spec 不会改变它。因为每个被跟踪的流都带有自己的 epoch,前端程序可以区分一次应用之前和之后选路的流,例如结束所有由较旧 plane 选路的流:
let current = supervisor.epoch();let killed = supervisor .tracker() .kill_where(|flow| flow.plane_epoch() < current);已打开的流保留什么
Section titled “已打开的流保留什么”路由变更从不迁移已打开的 stream:它留在打开时所在的目标上。已打开的 stream 携带其出站打开的 link,并在 Guarded 内携带该目标 closed token 的克隆;正是这个 token 让之后的排空能够触及它。stream 不会再读取 cell,它的 connect 取得的 load() guard 在 connect 返回时就已释放。UDP 关联则自己持有 cell,每次发送都读取它。
| 流打开之后的变化 | TCP 流,或一个 mux 子流 | UDP 关联 |
|---|---|---|
| 路由变更把它的目的地送往别处 | 留在打开时所在的目标上(a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow)。同一个 mux 载体的下一个子流走新路由(a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session)。 |
下一个数据包在新 plane 上选路,按新路由去往目的地(同一个 UDP 测试)。 |
| 它的出站被重建或移除 | actor 对旧版本的 token 应用该版本的排空策略(排空)。 | actor 以同样方式对现有子链路的 token 应用排空策略;路由现在指向新版本的数据包会在新版本上打开一条子链路,因为两者的 id 不同。 |
| 它的负载均衡器被重建 | 无影响:流受其成员的 token 保护,而不是负载均衡器的(a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers)。 |
后续数据包在新的负载均衡器上解析。 |
| 某个负载均衡器成员变为不可用 | 无影响:选择发生在流打开时。 | 后续数据包可能解析到另一个成员。 |
| 某个用户的限速改变 | 在流的下一次读或写时生效(流跟踪、统计与限速)。 | 同样生效,按子链路计。 |
重建或移除出站的应用会在 PublishPlane 之后为旧版本规划 Step::Drain { outbound, policy };默认是 DrainPolicy::Keep。在提交中,当 actor 之前的 targets map 恰好持有那个版本时,drain(supervisor/src/supervisor.rs)为旧目标运行。排空策略只决定 actor 如何处理该版本的 closed token:
DrainPolicy |
该版本的 token |
|---|---|
Keep |
actor 永不取消它,所以 Guarded 永远不会让这些流失败。新 plane 不再指向这个版本,所以不会有新流在它上面打开。 |
Close |
由 target.close_flows() 立即取消。 |
CloseAfter(grace) |
在 grace 之后由 close_flows() 取消,执行者是 supervisor TaskTracker 上的一个任务,它在 tokio::time::sleep(grace) 和根 token 之间 select;如果根 token 先被取消,任务不做取消就结束。 |
每个出站的策略如何确定,以及应用报告列出什么,见规划与应用变更。
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close 用一个 freedom 出站测试了前两种:在默认策略下,版本 1 上的流在应用之后仍能回显;在 Close 下做第二次变更之后,版本 2 上的流被关闭,版本 1 上的流不会被再次排空,新流使用版本 3。shutdown_ends_pending_grace_timers 验证一个一小时的 CloseAfter 既不会在关停之后继续存在,也不会拖延关停。
只有出站版本会被排空。supervisor 从不取消负载均衡器目标自己的 token,也不需要取消:经过负载均衡器的流受其成员的 token 保护,所以对该成员版本的 Close 或 CloseAfter 排空能够触及它。DNS 目标也从不被排空;没有任何 Drain 步骤指向它。
| 不变量 | 由谁保证 | 由哪些测试验证 |
|---|---|---|
| 路由和目标一起发布,规则永远不会指向 plane 中没有的 tag。 | 用同一次应用的路由和目标构建一个 Plane,并用一次 store 发布;校验拒绝未知的路由 tag;对槽位数量的 debug_assert_eq!。 |
a_refused_spec_changes_nothing(未知的默认路由被拒绝,epoch 不变);every_rule_naming_a_tag_shares_that_tags_slot |
| 每个 tag 一个槽位,decision 仍然指明它的规则。 | CompiledRoutes::new 中的 slot_of;每条规则一个 Arc<Decision>。 |
every_rule_naming_a_tag_shares_that_tags_slot |
首个匹配的规则胜出;没有匹配时走默认路由,rule: None。 |
RouteTable::pick。 |
every_rule_naming_a_tag_shares_that_tags_slot;路由模型自己的测试 |
| TCP 流只选路一次;路由变更作用于下一个流,而不是已建立的流。 | connect 只读取一次 cell,打开的 stream 不再查询它。 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
| 每个 mux 子流在打开时选路。 | 运行时对每个 Effect::Open 调用 connect。 |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
| UDP 关联在每次发送时按当前 plane 选路。 | FanOutLink::poll_send_to 每次调用都读取 cell。 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
| 存在 DNS 服务时,端口 53 经 TCP 和 UDP 都到达该服务,而服务自己的查询永远不会被拦截。 | Plane::route 与 Plane::route_upstream 的区分。 |
port_53_is_intercepted_except_for_the_services_own_queries、without_a_dns_service_port_53_follows_the_route_table |
| 第一次提交之前,没有任何流量离开主机。 | Plane::empty。 |
the_empty_plane_drops_everything |
| 入站 tag、源地址和网络类型都会到达路由器。 | route_target。 |
inbound_tag_selects_the_route、source_cidr_matches_the_client_address、network_separates_tcp_from_udp(app/tests/integration/e2e_route_context.rs);supervisor/tests/unit/router.rs 中的四个源地址测试 |
| 嗅探到的名称参与匹配,拨号的是请求的目的地。 | route_target 设置 sniffed_domain;没有任何东西改写 flow.destination。 |
an_http_host_routes_an_ip_addressed_flow、a_tls_sni_routes_an_ip_addressed_flow、a_flow_whose_sniffed_host_does_not_match_is_blocked、an_unsniffable_payload_falls_through_to_the_default、turning_sniffing_off_stops_the_domain_rule_matching(app/tests/integration/e2e_sniff.rs) |
| 经负载均衡器的流记在承载它的成员名下,并随该成员一起排空。 | 在 FlowMeta::of 之前调用 resolve(None);connect_stream 取解析后目标的 token。 |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers |
| 关闭型排空让流的每个方向都失败,并唤醒挂起的接收。 | 除 poll_shutdown 之外,每个受保护方法中的 Guarded::poll_closed。 |
closing_a_targets_flows_aborts_its_open_streams、closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive |
| 旧版本只在新 plane 发布之后才被排空。 | 规划器把 Drain 排在 PublishPlane 之后;提交在写入之后才遍历这些步骤。 |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
| 流先准入再选路。 | Session::admit 是 connect 的第一条语句。 |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy(supervisor/tests/unit/session.rs)覆盖准入本身 |
| 流只在其出站打开之后才注册,并带有规则、解析后的出站和 epoch。 | meter_stream 在 opened.await? 之后运行。 |
成功后才注册:没有专门的测试。元数据:mux_sub_flows_and_their_session_count_exactly_what_moved、a_udp_association_counts_each_sub_link_and_charges_its_session(supervisor/tests/tracking.rs) |
一个 tag 的版本号从不重复,所以以 OutboundId 为 key 永远不会把正在排空的构建与其后继混淆。 |
next_version;被移除 tag 的版本号保存在 RunningState 中。 |
each_rebuild_of_a_tag_takes_the_next_version、a_tag_added_back_takes_a_version_it_never_had |
| 被拒绝的应用让 plane 保持原样。 | 路由和目标在准备阶段构建;只有提交会写入。 | a_refused_spec_changes_nothing |
失败路径与取消
Section titled “失败路径与取消”| 位置 | 错误 | 后续 |
|---|---|---|
Session::admit:会话的 token 已被取消,或会话已不在注册表中 |
PermissionDenied,the session was closed |
future 立即 resolve 为该错误。运行时把该 key 的 Event::ConnectFailed 交给协议核心;SOCKS 驱动器从它的 await 拿到错误(SOCKS)。 |
Session::admit:第一个流出示了已被吊销的 principal |
PermissionDenied,the user was removed before the session opened a flow |
会话的 token 也被取消,整个会话随之结束。 |
| 解析得到的目标是负载均衡器 | balancer member <tag>@v<n> is itself a balancer |
只要校验成立就不会出现。 |
| 出站拨号失败 | 出站的错误 | 在 meter_stream 之前返回;不注册任何流。错误文本见出站、UDP 扇出与负载均衡器。 |
| 排空关闭了流所在的版本 | ConnectionAborted,the outbound this flow was opened on was drained |
由流的下一次读、写、flush、发送或接收返回。在运行时之下,协议核心收到 Event::OutboundError 并关闭这个流。以这种方式失败的 UDP 子链路会从其关联中丢弃。 |
| tracker 强制关闭流 | ConnectionAborted,the flow was killed |
由 Metered 产生,在 Guarded 之外(流跟踪、统计与限速)。 |
没有任何东西再持有 plane cell 时的 RoutedDialer |
NotConnected,dns: the supervisor this resolver dials through is gone |
解析器的查询失败。 |
准备阶段 compile_routes 失败 |
route 的 ApplyError::Build |
应用被拒绝;plane 不会重新发布。 |
| 规则指向未知 tag | ApplyError::UnknownReference,route references unknown outbound <tag> |
在准备之前被校验拒绝。 |
| 位置 | 级别与 target | 文本 |
|---|---|---|
FanOutLink,子链路打开失败 |
debug,etemenanki_supervisor::topology::outbound::udp_fanout |
udp fan-out: opening <tag>@v<n> failed: <error> |
FanOutLink,在子链路上发送失败 |
debug,同上 |
udp fan-out: sending on <tag>@v<n> failed: <error> |
FanOutLink,在子链路上接收失败 |
debug,同上 |
udp fan-out: sub-link <tag>@v<n> ended: <error> |
| 负载均衡器探测发现成员健康状态变化 | info,etemenanki_supervisor::topology::balancer |
balancer member <tag> is now up 或 balancer member <tag> is now down |
Actor::commit 结束时 |
debug,etemenanki_supervisor::supervisor |
applied: <n> built, <n> reused, <n> swapped, <n> drained |
build_geo_data |
warn,etemenanki_environment::routing |
skipping unknown geosite domain type <n>、skipping invalid geosite regex "<pattern>": <error> |
connect、plane、Guarded 和 RoutedDialer 不记录任何日志。
connect 和 plane 不 spawn 任何东西。调用方拥有返回的 future(运行时把它保存在该 key 的 LinkState::Connecting 中),drop 它就会取消拨号。
发布 plane 的提交会启动负载均衡器的探测任务,每个新负载均衡器的每个成员一个;CloseAfter 排空会启动一个定时器任务。两者都运行在 supervisor 的 TaskTracker 上。当负载均衡器的探测 token(根 token 的子 token)被取消时,探测任务结束:发生在不再复用该负载均衡器的那次发布时,或关停时。CloseAfter 任务在宽限期结束后或根 token 被取消时结束。探测循环见出站、UDP 扇出与负载均衡器。
| 项目 | 位置 | 值 |
|---|---|---|
| 被拦截的端口 | supervisor/src/topology/plane.rs → DNS_PORT |
53,TCP 和 UDP 都拦截 |
| 每张表的规则数 | supervisor/src/topology/router.rs → CompiledRoutes::new |
少于 2^32;RuleId 是 u32,编译更多规则会以 fewer than 2^32 rules panic |
| 每条规则内联存放的匹配条件 | environment/src/routing.rs → RouteItem |
3 个(SmallVec<[RouteMatch; 3]>);更多的存放在堆上,没有上限 |
| cell 读取 | supervisor/src/connector.rs、supervisor/src/topology/outbound/udp_fanout.rs、supervisor/src/topology/outbound/dns.rs |
每个 TCP 流(包括 mux TCP 子流)、每次 poll_send_to 调用、每条解析器连接各一次 ArcSwap::load |
| epoch 与版本号 | supervisor/src/supervisor.rs、supervisor/src/topology/spec_plan/plan.rs |
u64 计数器,每次发布、每次重建一个 tag 各加一 |
| 每个 UDP 关联的子链路数 | supervisor/src/topology/outbound/udp_fanout.rs → MAX_SUBS |
64(出站、UDP 扇出与负载均衡器) |
一次路由决策的代价是按顺序遍历规则直到首个匹配,每次决策把流的域名转为小写一次;各匹配条件的开销见路由模型页面。
| 层级 | 文件 | 覆盖内容 |
|---|---|---|
| 单元,plane | supervisor/tests/unit/plane.rs |
port_53_is_intercepted_except_for_the_services_own_queries、without_a_dns_service_port_53_follows_the_route_table、the_empty_plane_drops_everything、closing_a_targets_flows_aborts_its_open_streams、closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive、a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers。目标都是 blackhole,所以不涉及网络。 |
| 单元,路由器 | supervisor/tests/unit/router.rs |
the_circuits_own_source_wins_over_the_listeners、without_one_the_listeners_source_is_used、with_neither_there_is_no_source_to_match_on、a_circuit_source_works_with_no_listener_source、every_rule_naming_a_tag_shares_that_tags_slot。 |
| 单元,计划 | supervisor/tests/unit/plan.rs |
何时规划 PublishPlane,例如 the_first_plan_binds_and_builds_everything_then_publishes、an_unchanged_spec_reuses_everything_and_publishes_nothing 和 removing_a_balancer_alone_republishes_the_plane;Drain 跟在它之后(a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish);版本编号(each_rebuild_of_a_tag_takes_the_next_version、a_tag_added_back_takes_a_version_it_never_had)。其余测试见规划与应用变更。 |
| 热切换 | supervisor/tests/hot_swap.rs |
在连接存活时对真实 socket 重新配置:a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow、a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session、a_changed_outbound_keeps_its_old_flows_unless_drained_with_close、a_refused_spec_changes_nothing、an_established_connection_survives_an_apply_that_keeps_its_inbound、shutdown_ends_pending_grace_timers。 |
| 跟踪 | supervisor/tests/tracking.rs |
connector 登记的元数据:每个 mux 子流和每条 UDP 子链路的出站 tag 与 rule()。 |
| app 端到端 | app/tests/integration/e2e_route_context.rs |
inbound_tag_selects_the_route、source_cidr_matches_the_client_address、network_separates_tcp_from_udp,每个测试所用的配置中,该匹配条件都是流量与 blackhole 之间唯一的屏障。 |
| app 端到端 | app/tests/integration/e2e_sniff.rs |
位于 Xray 客户端之后的一个 VLESS 入站:HTTP Host 和 TLS SNI 为以 IP 寻址的流选路,不匹配的名称和无法嗅探的载荷落到默认路由,sniffing = false 让匹配失效。被阻断的流永远不会应答,所以往返上的 10 秒超时是该测试的「已丢弃」信号,而不是失败。 |