supervisor 概览
源码文件:54 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/Cargo.tomlEtemenanki/supervisor/src/lib.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/policy.rsEtemenanki/supervisor/src/entity/mod.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/user.rsEtemenanki/supervisor/src/entity/usage.rsEtemenanki/supervisor/src/build/mod.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/build/users.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/build/outbound.rsEtemenanki/supervisor/src/topology/mod.rsEtemenanki/supervisor/src/topology/spec_plan/mod.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/router.rsEtemenanki/supervisor/src/topology/inbound/mod.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/system/mod.rsEtemenanki/supervisor/src/system/listener.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/sampler.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/socket_policy.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/supervisor/tests/support/mod.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/validate.rsEtemenanki/supervisor/benches/tracking.rsEtemenanki/environment/src/dial/socket.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/lower.rsEtemenanki/app/src/api.rsEtemenanki/ffi/src/proxy.rsEtemenanki/ffi/src/platform.rsEtemenanki/ffi/src/client.rsEtemenanki/webclient/src/tracking.rsEtemenanki/app/tests/unit/subscribe.rsEtemenanki/app/tests/integration/e2e_balancer.rsEtemenanki/app/tests/integration/e2e_dns.rskatana/src/manager/node.rskatana/src/lower/mod.rskatana/src/api/mod.rs
etemenanki-supervisor 是每个 Etemenanki 程序都运行在其上的运行时。etemenanki-app、移动端库 etemenanki-ffi 和 katana 各自把自己的输入(一个 TOML 文件及其订阅文件、手机 app 传入的字符串、面板的应答)降为一个带类型的 Spec<U>,再交给一个 Supervisor。supervisor(监管器)校验这个 spec(期望状态),对照正在运行的内容规划一次调和,并分两个阶段应用它:先是可能失败、但不改动任何正在运行之物的准备(prepare),然后是不会失败的提交(commit)。调和会保留没有变化的部分:会话不属于接受它的监听器,未变的绑定保留它的 socket,未变的出站保留它的状态。这个 crate 对配置文件、配置格式以及重载如何触发一无所知。
本页是 supervisor 部分的入口。内容包括:supervisor 与其前端程序之间的分工、模块地图、公共 API、每次变更都要经过的 actor 以及一次应用(apply)的数据流、取消 token 与任务树、socket 策略、关停,以及为正在运行之物命名的 id。spec 的各个类型见 spec:期望状态,语义规则见校验与应用错误,计划、准备和提交的每一步见规划并应用变更。
supervisor 负责:
- 按每一条语义规则校验 spec,只要有一条不满足就整体拒绝它(这些规则只归
build/validate.rs所有); - 比较正在运行的 spec 与期望的 spec,逐个资源地规划是保留、构建、替换、停止、关闭还是排空;
- 构建出站、负载均衡器、路由表、DNS 解析器和 DNS 服务、入站 handler(处理器)以及用户表;
- 绑定监听器:TCP、Hysteria 2 用的 UDP、Unix socket,以及它自己创建或以描述符形式交给它的 TUN 设备;
- 运行 accept 循环,把每个被接受的连接登记为一个会话,并驱动每个连接的协议运行时;
- 通过当时的 plane(数据平面)为每个新流、mux 子流和 UDP 包选路;
- 为出站的每次构建保留一个版本,并把被替换或被移除的版本交给它的排空策略;
- 按入站准入用户,按移除策略撤销被移除的用户,并执行按用户的限速;
- 跟踪每个存活的流,采样速率,并维护按用户的用量账本;
- 关停时停止以上全部。
它不负责:
- 解析任何格式,或读取配置文件。证书、密钥和用户列表以字节和带类型的值的形式随 spec 到达。它读取的文件只有
RouteSpec指名的 geo 数据文件,以及Split形态的 DNS spec 设置了use_hosts时的系统 hosts 文件。 - 决定何时重载。文件监视、面板轮询或 API 调用都是前端程序的事;supervisor 只应答
apply。 - 记录配置错误日志。它返回一个
ApplyError,由前端程序决定如何展示。
一次应用保留什么
Section titled “一次应用保留什么”crate 文档(supervisor/src/lib.rs)列出了一次应用会保留的四样东西:
| 内容 | 如何保留 | 所属页面 |
|---|---|---|
| 连接 | 每个被接受的连接都是一个会话,处在它自己的取消 token 之下;这个 token 是 supervisor 根 token 的子 token,而不是其监听器的子 token。 | 用户、principal 与会话 |
| 监听器 | 监听器以其 BindSpec 为 key。未变的绑定保留它的 socket 或设备,只替换服务新连接的 handler。 |
监听器与服务循环 |
| 路由 | 路由和出站组成一个 Plane,放在一个原子 cell 之后。每个新流、mux 子流和 UDP 包都读取当前的 plane;路由变更不会为已经打开的流重新选路。 |
plane:为每个流选路 |
| 有状态的出站 | spec 未变、且解析器也未变的出站通过 Arc 沿用,保留它的 QUIC 连接、隧道或负载均衡器健康状态。DnsSpec 变化会重建除 blackhole 之外的每个出站(见规划并应用变更)。 |
出站、UDP fan-out 与负载均衡器 |
用户是一种独立的资源。它们通过 set_users、upsert_user 和 remove_user 变更,这些调用只重建准入了所变更用户集的那些入站的用户表。一个用户的限速由其所有流共享,并在存活的流下一次读写时生效,不会结束该用户的会话。
以下情况会结束存活的连接:
| 原因 | 由谁决定 | 默认 |
|---|---|---|
| 用户不再被某个入站准入:被移出其用户集,或其用于该入站的凭据变了 | UserRemovalPolicy,按入站或在 supervisor 全局设置 |
Close |
| 某个出站版本被替换或被移除 | DrainPolicy,按出站或在 supervisor 全局设置 |
Keep:不采取任何动作 |
| 某个入站从 spec 中被移除,或其 tag 被改名 | 总是如此:其会话被关闭 | 无 |
规划器标记为 Step::Disrupt 的变更:TUN 入站的设备或设置,或 Hysteria 2 监听器的混淆 |
需要 ApplyOptions::allow_disruptive;etemenanki-app 会拒绝它并提示重启 |
拒绝 |
close(selector) |
调用方 | 无 |
Tracker::kill(id) 或 Tracker::kill_where(select):一个流,或一个谓词选出的那些流 |
调用方:FFI 的 close_connection、REST API(webclient/src/tracking.rs) |
无 |
shutdown(grace) |
调用方;存活的连接先获得最多 grace 的时间 |
无 |
supervisor/src/policy.rs 定义了这两种策略及其 supervisor 全局默认值:
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub enum UserRemovalPolicy { Keep, #[default] Close, CloseAfter(Duration),}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub enum DrainPolicy { #[default] Keep, Close, CloseAfter(Duration),}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct Policies { pub user_removal: UserRemovalPolicy, pub drain: DrainPolicy,}UserRemovalPolicy作用于绑定到已撤销 principal(身份主体)的每个存活会话(entity/session.rs→enforce):Keep不动会话的 token,Close立即取消它,CloseAfter(grace)在grace之后取消它,除非会话先结束。DrainPolicy作用于旧的出站版本(supervisor.rs中的drain):Keep对旧版本的流不采取任何动作,Close立即取消该版本的 closed token(Target::close_flows),CloseAfter(grace)则在grace之后这样做。- 两个默认值有意不同。移除用户默认
Close,因为撤销访问权限本来就是要让它生效。替换出站默认Keep,因为 TCP 流无法在中途迁移到后继版本,所以默认情况下排空不会关闭旧版本的流。 Spec::policies设置 supervisor 全局的值。InboundSpec::user_removal为单个入站覆盖移除策略(supervisor.rs中的removal_policy回退到spec.policies.user_removal),OutboundSpec::drain为单个出站覆盖排空策略(由规划器解析,见规划并应用变更)。- 两种
CloseAfter定时器都运行在 supervisor 的TaskTracker上,所以关停会等待它们。移除定时器在其会话的 token 上 select,排空定时器在根 token 上 select,而关停会取消根 token(会话 token 是它的子 token),所以两种定时器都会立即结束,而不会等满宽限期(见任务树)。
前端程序负责什么
Section titled “前端程序负责什么”| 事项 | 前端程序 | supervisor |
|---|---|---|
| 输入格式 | 解析自己的格式:TOML、订阅文件、面板的 JSON、FFI 参数 | 只接受 Spec<U> |
| 解析与 lowering(降为 spec) | 报告语法错误和未知键;解析自己的字符串(UUID、加密方法名、地址、密钥、CIDR);应用自己 schema 的默认值;为自己的用户命名;拒绝 spec 无法表达的键组合 | 完全看不到这些 |
| 语义(重复的 tag、未知的引用、各个带类型的部分如何组合在一起) | 不重复检查 | build/validate.rs,唯一的所有者 |
| 文件 | 读取证书、密钥和用户列表;传递字节 | 只读取 geo 数据和 hosts 文件 |
| 用户身份 | 选择 key 类型 U |
为每个被准入的 U 分配一个内部 UserKey |
| 何时变更 | 监视文件、轮询面板、提供 API | 应答 apply、update 和用户编辑 |
| 中断性变更 | 选择 ApplyOptions |
没有 allow_disruptive 时拒绝 |
| 主机 socket 策略 | 把 SocketOptions 传给 builder |
应用到每个出站 socket |
| 用量 | 轮询 take_usage 或安装一个 UsageSink |
按用户统计线路字节 |
| 日志 | 配置 tracing,并用自己的文字包装 ApplyError |
通过 tracing 以自己的 target 记录日志 |
actor 及其启动的监听器会记录以下日志行。accept 循环和连接记录的日志见监听器与服务循环。
| 级别 | target | 日志行 | 时机 |
|---|---|---|---|
| debug | etemenanki_supervisor::supervisor |
applied: {} built, {} reused, {} swapped, {} drained |
每次提交结束时,带上报告中这四个字段的长度 |
| info | etemenanki_supervisor::system::listener |
inbound {tag} listening on {bind} |
Listener::start,提交启动的每个监听器各一行 |
| info | etemenanki_supervisor::system::listener |
inbound {tag} owns tun device {name} |
创建了 TUN 设备的绑定,带设备名 |
| info | etemenanki_supervisor::system::listener |
inbound {tag} serves a supplied tun device |
接管了外部提供的描述符的绑定 |
| debug | etemenanki_supervisor::system::listener |
could not remove {path}: {e} |
Unix 监听器被 drop,而它的 socket 文件(仍是它创建的那个)无法删除(SocketFile::drop) |
supervisor/Cargo.toml 声明 etemenanki-supervisor 的版本为 0.3.1,只发布到一个私有 Cargo registry。它依赖 etemenanki-concepts 2.0.0、etemenanki-environment 3.0.0 和 etemenanki-protocols 4.0.1,并且始终打开 hysteria 和 tun 两个 feature,所以每个 supervisor 都能服务 Hysteria 2 和 TUN。在其余依赖中,manifest 自己的注释说明了其中三个的用途:scc 是存活流的注册表,quinn 只用来指称 Hysteria 2 监听器被替换成的 TLS 材料,serde 让 UsageDelta 可以原样送到面板(“No config types: this crate never parses a config format”)。
开发依赖有:带 test-util 的 tokio(以便在暂停的时钟下测试宽限期)、用于测试证书的 openssl、用于用量账本属性测试的 proptest、用于 tracking 基准测试的 criterion,以及 socket2,它提供 socket 策略 hook 拿到的 socket 类型。
文件夹supervisor/
- Cargo.toml
文件夹src/
- lib.rs crate 文档、模块与重新导出
- supervisor.rs
Supervisor、SupervisorBuilder、check以及 actor - policy.rs
UserRemovalPolicy、DrainPolicy、Policies - serve.rs accept 循环与服务单个连接(crate 私有)
- connector.rs
AppConnector,每个入站拨号都经过的 connector(连接器) 文件夹build/ 从已校验的 spec 构建运行所需的一切
- mod.rs
- apply.rs
ApplyReport、ApplyError - validate.rs spec 的每一条语义规则
- dns.rs 解析器与 DNS 服务
- inbound.rs 入站 handler 及其用户表
- outbound.rs 每种协议一个出站 handler
- route.rs 路由表及其 geo 数据
- users.rs 准入、principal 与用户 key
文件夹entity/
- mod.rs
- id.rs
OutboundId、Resource、SessionId、UserKey、FlowId、RuleId - user.rs
UserId、UserName、UserSpec、UserSet、凭据 - session.rs 会话注册表
- usage.rs 线路计数器、用量账本、
UsageSink
文件夹system/
- mod.rs
- listener.rs 绑定
BindSpec,以及服务它的任务
文件夹topology/
- mod.rs
文件夹spec_plan/
- mod.rs
Spec、ObfsSpec、Secret - inbound.rs 入站 spec
- outbound.rs 出站 spec
- transport.rs stream 形态与 TLS spec
- route.rs
RouteSpec、BalancerSpec - dns.rs
DnsSpec - plan.rs
RunningState、Plan、Step、plan
- mod.rs
文件夹inbound/
- mod.rs
BindSpec与入站 handler
- mod.rs
文件夹outbound/
- mod.rs
Outbound,每种类型一个变体 - proxy.rs 代理协议客户端
- freedom.rs 直连 TCP 与 UDP
- udp_fanout.rs
FanOutLink,按包为 UDP 关联选路 - dns.rs 作为出站的 DNS 服务,
RoutedDialer
- mod.rs
- plane.rs
Plane、PlaneCell、Target - router.rs 编译后的路由表
- balancer.rs 负载均衡器及其健康探测
- flow.rs
Flow、Principal、FlowContext
文件夹track/
- mod.rs
Tracker、FlowEntry、FlowEvent - metered.rs
Metered、MeteredDatagram - pace.rs 按用户的令牌桶
- sampler.rs 采样器与
StatsSnapshot
- mod.rs
文件夹tests/
- hot_swap.rs 针对存活连接的应用
- socket_policy.rs socket 策略 hook
- tracking.rs 流、会话与用量
文件夹support/
- mod.rs 共享的 spec、回显服务器、协议客户端
文件夹unit/ 单元测试,编译进 crate
- …
文件夹benches/
- tracking.rs
lib.rs 把 build、connector、entity、policy、system、topology 和 track 设为公开,serve 保持 crate 私有,supervisor 保持私有,并在 crate 根重新导出各个入口:
pub use build::apply;pub use supervisor::{ApplyOptions, Selector, SessionInfo, Supervisor, SupervisorBuilder, check};在 build 内部,只有 apply 和 validate 是公开的;各个构建器(dns、inbound、outbound、route、users)是 crate 私有的。system 没有公开项:system::listener 是 crate 私有的。单元测试位于 tests/unit/,通过 #[path] 属性编译进它们所测试的模块(plan.rs、validate.rs、session.rs、plane.rs、router.rs、balancer.rs、track.rs、usage.rs、serve.rs 和 dns_outbound.rs)。
| 模块 | 所在页面 |
|---|---|
topology::spec_plan(类型)、policy、entity::user(类型) |
spec:期望状态 |
build::validate、build::apply |
校验与应用错误 |
topology::spec_plan::plan、actor 的准备与提交 |
规划并应用变更 |
serve、system::listener、topology::inbound、build::inbound |
监听器与服务循环 |
entity::user、entity::session、build::users |
用户、principal 与会话 |
connector、topology::plane、topology::router、topology::flow、build::route |
plane:为每个流选路 |
topology::outbound、topology::balancer、build::outbound |
出站、UDP fan-out 与负载均衡器 |
build::dns、topology::outbound::dns、topology::spec_plan 中的 DNS 类型 |
名称解析与 DNS 服务 |
track |
流跟踪、统计与限速 |
entity::usage |
按用户的用量计费 |
supervisor(句柄、builder、actor)、entity::id |
本页 |
Supervisor<U>
Section titled “Supervisor<U>”#[derive(Clone)]pub struct Supervisor<U: UserId> { commands: mpsc::Sender<Command<U>>, plane: PlaneCell, tracker: Tracker, usage: Arc<UsageBook<U>>,}
impl<U: UserId> Supervisor<U> { pub fn builder() -> SupervisorBuilder<U>; pub async fn start(spec: Spec<U>) -> Result<(Self, ApplyReport), ApplyError>;
pub async fn apply(&self, spec: Spec<U>) -> Result<ApplyReport, ApplyError>; pub async fn apply_with(&self, spec: Spec<U>, options: ApplyOptions) -> Result<ApplyReport, ApplyError>; pub async fn update(&self, edit: impl FnOnce(&mut Spec<U>) + Send + 'static) -> Result<ApplyReport, ApplyError>; pub async fn update_with( &self, edit: impl FnOnce(&mut Spec<U>) + Send + 'static, options: ApplyOptions, ) -> Result<ApplyReport, ApplyError>;
pub async fn set_users(&self, set: &str, users: BTreeMap<U, UserSpec>) -> Result<(), ApplyError>; pub async fn upsert_user(&self, set: &str, id: U, user: UserSpec) -> Result<(), ApplyError>; pub async fn remove_user(&self, set: &str, id: U) -> Result<bool, ApplyError>;
pub async fn close(&self, selector: Selector<U>) -> usize; pub async fn sessions(&self) -> Vec<SessionInfo<U>>;
pub fn epoch(&self) -> u64; pub fn tracker(&self) -> &Tracker; pub fn take_usage(&self) -> Vec<UsageDelta<U>>; pub fn usage_snapshot(&self) -> Vec<UsageTotal<U>>;
pub async fn shutdown(&self, grace: Duration);}Supervisor 是一个句柄。克隆它得到的是指向同一个运行中 supervisor 的另一个句柄:命令发送端、plane cell、tracker 和用量簿都是共享的。最后一个句柄在没有调用 shutdown 的情况下被 drop 时,actor 会看到它的命令通道关闭,并立即停止一切(见关停)。
| 方法 | 是否经过 actor | 作用 |
|---|---|---|
builder |
否 | 一个全部取默认值的 SupervisorBuilder。 |
start |
自己运行第一次应用 | Self::builder().start(spec)。 |
apply、apply_with |
Command::Apply |
调和到 spec。apply 传入 ApplyOptions::default()。 |
update、update_with |
Command::Update |
调和到经 edit 修改后的运行中 spec。 |
set_users |
Command::Users |
替换用户集 set 中的全部用户。 |
upsert_user |
Command::Users |
向 set 添加一个用户,或替换其 UserSpec。 |
remove_user |
Command::Users |
从 set 移除一个用户;若该用户原本在其中,返回 true。 |
close |
Command::Close |
立即关闭 selector 指定的会话;返回本次调用关闭的数量。 |
sessions |
Command::Sessions |
每个存活会话,按 id 排列。 |
epoch |
否 | 已发布的 plane 数量:self.plane.load().epoch()。 |
tracker |
否 | Tracker:存活的流与会话、流事件、采样得到的速率。 |
take_usage |
否 | 自上一次取出以来每个用户的线路字节。 |
usage_snapshot |
否 | 自启动以来每个用户的线路字节,无论是否已被取出。 |
shutdown |
Command::Shutdown |
停止接受连接,给存活连接最多 grace 的时间,关闭剩下的连接,并等待它们结束。 |
apply 替换整个 spec,包括用户集。update 从 actor 持有的 spec 出发,其中包含自上一次应用以来提交的每一次用户编辑,所以一个逐个编辑用户、之后又修改某个监听器的前端程序可以使用 update,而不必重新发送它的用户。
take_usage 为每个用户最多返回一个 UsageDelta,没有搬运任何字节的用户则没有。它对每个活跃用户访问一个账户。代价是滞后:存活会话的字节在每个采样器 tick 时统计,所以它自上一个 tick 以来搬运的字节会出现在之后的某次取出中;已结束的会话则全部计入。用量 sink 取走的字节永远不会由 take_usage 返回,反之亦然。账本本身见按用户的用量计费。
SupervisorBuilder<U>
Section titled “SupervisorBuilder<U>”pub struct SupervisorBuilder<U: UserId> { options: ApplyOptions, sample_interval: Duration, usage_sink: Option<PushUsage<U>>, socket: SocketOptions,}
impl<U: UserId> Default for SupervisorBuilder<U>;
impl<U: UserId> SupervisorBuilder<U> { pub fn options(self, options: ApplyOptions) -> Self; pub fn sample_interval(self, interval: Duration) -> Self; pub fn usage_sink(self, sink: impl UsageSink<U>, interval: Duration) -> Self; pub fn socket_options(self, socket: SocketOptions) -> Self; pub async fn start(self, spec: Spec<U>) -> Result<(Supervisor<U>, ApplyReport), ApplyError>;}builder 保存的是比任何一个 spec 都长寿的设置:
| 设置 | 默认值 | 效果 |
|---|---|---|
options |
ApplyOptions::default() |
第一次应用的选项。在它之前没有任何东西在运行,所以第一次计划从不包含中断性步骤,这个选项对它不起作用。 |
sample_interval |
DEFAULT_SAMPLE_INTERVAL,1 s |
采样器多久发布一次 StatsSnapshot 并调和一次用量账本。传入 Duration::ZERO 时,start 以 the sample interval must not be zero panic。 |
usage_sink |
无 | 每个 interval 把每个用户的用量交给 sink 一次,并在关停时、最后一个会话结束之后交出最后一批。 |
socket_options |
SocketOptions::default() |
在 supervisor 整个生命周期内,每个出站 socket 打开时所遵循的策略(见 socket 策略)。 |
usage_sink 保存一个装箱的闭包(PushUsage<U>),start 把它变成一个运行 push_usage 的任务。sink 在这个任务上运行,所以慢的 sink 永远不会拖住连接:它只会推迟自己的下一批,而那一批随之覆盖更长的间隔。这个任务的工作方式如下:
- 第一批在
start之后一个interval到来(tokio::time::interval_at(now + interval, interval));sink 较慢时错过的 tick 不会集中补上(MissedTickBehavior::Delay)。 - 在每个 tick,以及
background_stop被取消时再一次,它在阻塞线程池上调和账本(spawn_blocking(move || ledger.reconcile())),所以这一批包含截至那一刻搬运的全部字节,而不是只到上一个采样器 tick。调和中的 panic 会用std::panic::resume_unwind在该任务上重新抛出。 - 然后它取出用量(
UsageBook::take),并且只在这一批非空时调用sink.report(&batch),正如UsageSink的约定所承诺的(“Never called with an empty batch”)。
批次如何形成见按用户的用量计费。
采样器(track/sampler.rs → Sampler::run)在 start 之后一个 sample_interval 第一次 tick(interval_at),并使用 MissedTickBehavior::Delay,所以因运行时繁忙而推迟的 tick 不会集中补上。每个 tick 都在阻塞线程池上运行(spawn_blocking,panic 用 resume_unwind 重新抛出),调和用量账本,并用 watch::Sender::send_replace 发布它的 StatsSnapshot。background_stop 被取消时,它运行最后一个 tick 然后返回。tick 本身见流跟踪、统计与限速。
start 依次执行:
- 断言
sample_interval不为零。 - 用 socket 策略构建一个
Actor(Actor::new),它会创建根 token、TaskTracker、用量账本、会话注册表、tracker 及其(尚未运行的)采样器,以及空 plane。 - 直接在调用方的任务上运行第一次应用:
actor.apply(spec, options).await?。被拒绝的 spec 在这里返回它的ApplyError。此时还没有 spawn 任何东西,被拒绝的那次准备所绑定的监听器,已在准备阶段返回错误时释放。 - spawn 采样器(
tokio::spawn(sampler.run(sample_interval, background_stop)))。 - 设置了 sink 时,在同一个
background_stoptoken 下 spawn 用量推送任务。 - 创建命令通道
mpsc::channel(16)。 - 用 actor 的 plane cell、tracker 和用量簿构建句柄。
- spawn actor 的循环
tokio::spawn(actor.run(rx)),并返回句柄和第一个ApplyReport。
第一次提交启动的监听器在第 3 步就开始接受连接,早于采样器运行。这没有害处:无论采样器何时第一次 tick,它都会统计各个流已经搬运的字节。
ApplyOptions
Section titled “ApplyOptions”#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct ApplyOptions { pub allow_disruptive: bool,}规划器把两类变更标记为 Step::Disrupt,因为它们都会结束存活的连接:对 TUN 入站的变更(其设备或设置),以及对 Hysteria 2 监听器混淆的变更。这样的变更需要 allow_disruptive。没有这个标志时,actor 在规划之后、准备任何东西之前拒绝整个 spec:
inbound <tag>: <reason>; this ends its live connections and needs allow_disruptiveactor 只报告计划顺序中的第一个中断(plan.disruptions().next()),所以即使有多个入站需要这个标志,ApplyError::Disruptive 也只指名一个入站。原因文字和各个步骤见规划并应用变更。
etemenanki-app 从不设置这个标志,所以它会拒绝这样的重载并提示重启,重启 app 即可应用变更。katana 总是设置它,这样来自面板的 Hysteria 2 混淆变更就能通过。a_hysteria2_obfuscation_change_needs_allow_disruptive 固定了 Hysteria 2 的情形:没有标志时被拒绝;有标志时,restarted 列出该入站,端口保持不变。
Selector<U>
Section titled “Selector<U>”#[derive(Debug, Clone, PartialEq, Eq)]pub enum Selector<U> { Session(SessionId), User(U), Inbound(CompactString), All,}| 选择器 | 关闭的内容 |
|---|---|
Session(id) |
该会话(如果它仍存活)。 |
User(id) |
绑定到该用户的每个会话,不论在哪个入站上。没有被任何入站准入过的用户 id 没有 key,什么也不关闭。第一个流还没有表明身份的会话尚不属于任何用户。 |
Inbound(tag) |
该入站的每个会话,无论是否绑定到用户。 |
All |
每个存活会话。 |
actor 把选择器映射为注册表的一个范围(Scope::One、Scope::User、Scope::Inbound、Scope::All),用户 id 会先被映射为它的 UserKey。计数的是本次调用取消的会话:已经关闭的会话不会被重复计数。
SessionInfo<U>
Section titled “SessionInfo<U>”#[derive(Debug, Clone, PartialEq, Eq)]pub struct SessionInfo<U> { pub id: SessionId, pub inbound: CompactString, pub source: Option<IpAddr>, pub user: Option<U>, pub up: u64, pub down: u64, pub started: Instant,}sessions 按会话 id 的顺序,从注册表的 SessionStats 构建这些值,并把每个 UserKey 转换回前端程序的 U。在会话的第一个流把它绑定到某个用户之前,user 为 None。对于从未认证的会话,以及不带用户而被准入、principal 为 Principal::anonymous() 的会话(例如在没有指名用户集的入站上),它始终为 None。up 和 down 是线路字节(见按用户的用量计费)。started 是一个 std::time::Instant。
Tracker::sessions 不经过 actor,以 SessionStats 的形式返回同一个注册表(带 UserKey 和用户标签,而不是 U)。REST API 走的就是这条路径。
pub async fn check<U: UserId>(spec: &Spec<U>) -> Result<(), ApplyError>;check 是对首次启动的一次演练(dry run)。它用 SocketOptions::default() 构建一个用完即弃的 Actor,从空的运行状态出发规划,并在关闭绑定的情况下运行准备阶段(prepare(&plan, spec, false)),然后 drop 它构建的一切。
check 会 |
check 不会 |
|---|---|
运行每一条校验规则(plan 首先调用 validate) |
绑定 TCP、UDP 或 Unix 监听器 |
| 构建 DNS 解析器,spec 要求时读取 hosts 文件 | 创建或接管 TUN 设备 |
| 构建每个出站,加载其 TLS 材料 | 把 handler 与已绑定的 socket 或设备配对(Pending::pair),因此不会打开 Hysteria 2 endpoint,也不会为运行时准备 TUN 设备 |
| 构建每个负载均衡器 | spawn 任务、启动采样器或拨号任何目标 |
| 编译路由表,读取其 geo 数据 | 检查中断性变更(没有任何东西在运行) |
| 准入用户并构建每个入站 handler,解析证书和私钥 | 使用主机的 socket 策略(它从来不需要) |
启动时会拒绝的,它都会拒绝,只有绑定以及把已绑定的 socket 或设备与其 handler 配对时所做的事除外:打开 Hysteria 2 endpoint(及其混淆),以及为运行时准备 TUN 设备。etemenanki-app 在两个地方通过 instance::check_bytes 调用它:--test 参数(instance::check,检查 -c 指定的文件),以及 REST API 的路由切换(app/src/api.rs → set_route),它在写入编辑后的配置之前先检查一遍。以下输出取自固定版本的二进制(已去掉时间戳和颜色):
$ etemenanki-app --test -c ok.tomlConfiguration OK.
$ etemenanki-app --test -c duplicate-tag.tomlERROR etemenanki_app: configuration invalid: duplicate outbound tag direct
$ etemenanki-app --test -c bad-ca.tomlERROR etemenanki_app: configuration invalid: building outbound proxy@v1 failed: error:04800064:PEM routines:PEM_read_bio_ex:bad base64 decode:../crypto/pem/pem_lib.c:968:failed: 之后的文字来自 OpenSSL,会随二进制链接的 OpenSSL 版本而变化。第二个错误来自校验。第三个来自准备阶段构建出站:app 读取了 CA 文件并把其字节放进 spec,而 supervisor 的出站构建器(与启动时用的是同一个)无法解析它们。configuration invalid: 是 app 的前缀;其余部分是 ApplyError 的 Display。
用户 id 类型 U
Section titled “用户 id 类型 U”Supervisor、它的 builder,以及每个指称用户的公共类型,都以前端程序的用户 key 为泛型参数:
pub trait UserId: Clone + Eq + Ord + Hash + Debug + Send + Sync + 'static { fn name(&self) -> UserName;}U 是 UserSet 的 key,并在每个 UsageDelta、UsageTotal、SessionInfo 和 Selector::User 中指称用户。name 给出日志和各协议用户标签所带的 UserName(一个 email 或用户名,从来不是机密)。数据平面从不接触 U:actor 为每个被准入的 id 分配一个 UserKey,连接携带的是这个 u64(见实体 id)。etemenanki-app 和 FFI 直接用 UserName 作为用户的 key;katana 用面板的数字 id。这个 trait 和用户集见用户、principal 与会话。
ApplyReport 与 ApplyError
Section titled “ApplyReport 与 ApplyError”两者都位于 supervisor/src/build/apply.rs。应用或 update 的应答是一个 ApplyReport 或一个 ApplyError;用户编辑的应答是 () 或一个 bool,或者一个 ApplyError。ApplyReport 列出这次应用对每个资源做了什么;期望 spec 中的每个资源都出现在 reused 或 built 中。提交阶段根据计划的步骤填写它:
| 字段 | 来源 |
|---|---|
reused |
每个 Step::Reuse |
built |
每个 Step::Build |
swapped |
每个 Step::SwapHandler |
drained |
每个 Step::Drain。只有当该 tag 正在运行的目标就是这个步骤指名的版本时(target.id() == outbound),排空策略才会执行,但无论如何都会列出该版本。 |
rebound |
针对某个入站的 Step::Bind,且该入站的 tag 在运行中的 spec 里有不同的绑定;针对新 tag 的 Step::Bind 只报告在 built 下 |
restarted |
每个 Step::Disrupt |
removed |
每个 Step::CloseSessions |
各字段的含义见规划并应用变更。
ApplyError 说明一次变更为什么没有执行,被拒绝的变更让正在运行的一切保持原样。各变体可能出现的位置:
| 变体 | Display |
产生位置 |
|---|---|---|
DuplicateTag |
duplicate {kind} tag {tag} |
plan → validate |
UnknownReference |
{from} references unknown {kind} {tag} |
plan → validate |
EmptyBalancer |
balancer {balancer} has no members |
plan → validate |
UnprobeableMember |
balancer {balancer}: outbound {member} has no upstream a TCP health probe can reach |
plan → validate |
Invalid |
{resource}: {reason} |
plan → validate;用户编辑的 validate_admission;指名不存在用户集的用户编辑 |
Disruptive |
inbound {inbound}: {reason}; this ends its live connections and needs allow_disruptive |
actor,在规划与准备之间 |
Build |
building {resource} failed: {source} |
准备阶段;用户编辑构建入站的用户表或使其适配监听器时 |
Bind |
inbound {inbound}: binding {bind} failed: {source} |
准备阶段,仅在绑定时 |
Stopped |
the supervisor has shut down |
句柄,当 actor 已不存在时;actor,在没有运行中 spec 时处理 update 或用户编辑(成功启动之后不可达) |
每条规则及其确切的原因文字见校验与应用错误。
一个任务拥有正在运行的一切。句柄通过通道与它通信,所以应用和用户变更会被串行化,而不需要跨 .await 持有锁。
type Reply<T> = oneshot::Sender<T>;type Edit<U> = Box<dyn FnOnce(&mut Spec<U>) + Send>;
enum Command<U: UserId> { Apply(Spec<U>, ApplyOptions, Reply<Result<ApplyReport, ApplyError>>), Update(Edit<U>, ApplyOptions, Reply<Result<ApplyReport, ApplyError>>), Users(CompactString, UserEdit<U>, Reply<Result<bool, ApplyError>>), Close(Selector<U>, Reply<usize>), Sessions(Reply<Vec<SessionInfo<U>>>), Shutdown(Duration, Reply<()>),}
enum UserEdit<U> { Set(BTreeMap<U, UserSpec>), Upsert(U, UserSpec), Remove(U),}每个经过 actor 的调用都使用同一个私有辅助函数 ask:
- 为应答创建一个
oneshot通道。 - 在有界的
mpsc通道上执行send(command).await。如果已有 16 个命令在排队,调用方就在这里等待空位。 - 等待应答。
第 2 步失败(接收端已不存在)或第 3 步失败(应答发送端被 drop)都会变成 ApplyError::Stopped。close 把它映射为 0,sessions 映射为空 Vec,shutdown 则忽略它。
Actor::run 每次接收一个命令,处理完之后才取下一个:
| 命令 | 处理者 | 内部等待 |
|---|---|---|
Apply(spec, options, reply) |
apply |
准备阶段创建 TUN 设备 |
Update(edit, options, reply) |
克隆运行中的 spec,在其上运行 edit,然后 apply;没有运行中的 spec 时应答 Stopped |
准备阶段创建 TUN 设备 |
Users(set, edit, reply) |
edit_users |
无 |
Close(selector, reply) |
close |
无 |
Sessions(reply) |
sessions |
无 |
Shutdown(grace, reply) |
shutdown,然后应答并返回 |
宽限期、被跟踪的任务、后台任务 |
commands.recv() 返回 None 时,说明所有句柄都在没有关停的情况下被 drop 了,循环会先运行 shutdown(Duration::ZERO) 再返回。
单任务带来的结果:
- 应用与用户编辑从不交错执行。不会有两次应用同时准备,所以两次准备永远不会争抢绑定同一个端口。
update的编辑在 actor 任务上、针对 actor 持有的 spec 运行,所以在读取运行中 spec 与应用编辑后的 spec 之间,不会插入任何其他变更。close和sessions要排在正在运行的应用之后。应用中唯一可能耗时的 await 位于准备阶段创建 TUN 设备之处(BindSpec::Tun(TunSource::Create(..))→tun::open(spec).await);接管外部提供的描述符、其他绑定以及每个构建器都在 actor 任务上同步运行。- actor 没有自己的锁。它的状态归该任务所有,通过
&mut self修改。它接触到的锁只在同步代码段内持有:会话注册表的 mutex、用量簿的RwLock和 tracker 的 pacer map(都是parking_lot锁)、流式监听器的watch通道(send_replace、borrow),以及 Hysteria 2 入站的 endpoint 锁,set_quic和set_obfs会获取它。
绕过 actor 的读取
Section titled “绕过 actor 的读取”必须保持低开销的读取从不排在应用之后:
| 读取 | 读的是 | 为什么不经过 actor 也安全 |
|---|---|---|
epoch() |
plane cell,ArcSwap::load |
plane 通过一次原子 store 发布。 |
tracker() |
共享的 Tracker |
流自行注册和注销;注册表是一个并发 map。 |
take_usage()、usage_snapshot() |
用量簿 | 账本有自己的锁:一个 mutex 保护它的账户 map(Ledger::inner,其中有 accounts 和 active),每个账户再各有一个 mutex。从 key 到 id 的名字表由 actor 在 RwLock 下追加。 |
这三种读取在关停之后仍然可用。
actor 的状态
Section titled “actor 的状态”struct Actor<U: UserId> { state: RunningState<U>, shared: Shared, root: CancellationToken, dns: Option<Dns>, dns_target: Option<Arc<Target>>, targets: HashMap<CompactString, Arc<Target>>, probes: HashMap<CompactString, CancellationToken>, routes: Option<Arc<CompiledRoutes>>, listeners: Vec<Listener>, admissions: HashMap<CompactString, Admissions<U>>, keys: UserKeys<U>, usage: Arc<UsageBook<U>>, sampler: Option<Sampler>, background: Vec<JoinHandle<()>>, background_stop: CancellationToken, epoch: u64, socket: SocketOptions,}| 字段 | 内容 |
|---|---|
state |
最后一次应用的 spec,以及每个出站和负载均衡器 tag 的最新版本,包括已移除的 tag(RunningState)。 |
shared |
每个 accept 循环共享的内容(serve.rs → Shared):plane cell、会话注册表、流的 Tracker 以及 TaskTracker。 |
root |
根取消 token。 |
dns、dns_target |
由运行中的 DnsSpec 构建的解析器,以及 supervisor 自己应答 DNS 时到达 DNS 服务所经的目标。 |
targets |
每个出站和负载均衡器 tag 的当前版本,类型为 Arc<Target>。 |
probes |
每个运行中负载均衡器的探测 token,按 tag 索引。 |
routes |
编译后的路由表。 |
listeners |
每个已绑定的监听器及其服务任务。 |
admissions |
每个入站准入了谁,按入站 tag 索引:用户 id → principal 与凭据。 |
keys |
为每个被准入的用户 id 分配的 UserKey(build/users.rs → UserKeys)。 |
usage |
与每个句柄共享的用量簿。 |
sampler |
采样器,直到 start spawn 它为止。 |
background、background_stop |
采样器和用量推送任务,以及停止它们的 token。 |
epoch |
迄今为止发布的 plane 数。 |
socket |
每个出站 socket 打开时遵循的 socket 策略。 |
用量簿是 actor 构建的状态中唯一一个由句柄直接读取的部分:
struct UsageBook<U> { ledger: Arc<Ledger>, names: RwLock<Vec<U>>,}names 按 key 的顺序列出每个分配过 key 的 id;actor 在提交新 key 时(commit_keys)追加它们的 id。take 和 totals 把每个账本 key 转换回 U。一个没有名字的 key 意味着某个会话在 key 发布之前就绑定了它,而已取出的字节无法放回,所以 UsageBook::name 会以 {key:?} was bound before it was published panic,而不是丢弃这个增量。
sequenceDiagram
participant FE as 前端程序
participant H as Supervisor 句柄
participant A as actor 任务
participant R as 正在运行的部分
FE->>H: apply_with(spec, options)
H->>A: 通道上的 Command::Apply
A->>A: plan(state, spec),先校验
alt spec 违反某条规则
A-->>H: Err(DuplicateTag、UnknownReference、EmptyBalancer、UnprobeableMember 或 Invalid)
else 有 Disrupt 步骤而没有 allow_disruptive
A-->>H: Err(Disruptive)
else 合法,且没有中断或已允许
A->>A: 准备:构建并绑定,正在运行的部分不变
alt 准备失败
A-->>H: Err(Build 或 Bind),准备构建的内容被 drop
else 准备成功
A->>R: 提交:key、限速、plane、撤销、监听器、排空
A-->>H: Ok(ApplyReport)
end
end
H-->>FE: 结果
-
规划(
topology/spec_plan/plan.rs→plan)校验 spec 并与RunningState比较。它是纯函数:不碰 socket、不读文件、不看时钟。 -
中断检查(
Actor::apply):没有allow_disruptive时,第一个Step::Disrupt会让 spec 以ApplyError::Disruptive被拒绝。 -
准备(
Actor::prepare)完成所有可能失败的工作,且不改动任何正在运行的东西:不启动监听器,不发布 plane,不改动用户表或 key。依次为:- DNS:除非计划中有
Step::Build(Resource::Dns),否则保留运行中的解析器。重建解析器(build_dns)时,如果 spec 让 supervisor 自己应答 DNS,也会构建 DNS 服务的目标Target::outbound(internal_id(self.epoch + 1), Outbound::Dns(service));否则保留运行中的dns_target。 - 出站:每个
Step::Build调用build_outbound,每个Step::Reuse沿用运行中的Arc<Target>。 - 负载均衡器:对每个
Step::Build,为每个成员出站创建一个带探测目标的Member,调用Balancer::new,并在计划的版本下创建一个Target::balancer;新的负载均衡器连同其probe_interval和probe_timeout被保存下来,留给提交阶段。Step::Reuse保留运行中的目标。 - 路由表:只有在
Step::Build(Resource::Route)时才重新编译(compile_routes)。 - 准入:对每个入站,在用户 key 的暂存副本上执行
admit,收集不再被准入的 principal 以及该入站的移除策略。 - 入站,逐个处理。计划要构建的入站会得到一个新的 handler(
build_handler)。在运行中的监听器上重建的 Hysteria 2 handler,如果max_circuits没变,会保留该监听器的 circuit 信号量(same_circuits,它把Listener::circuits()传给构建器)。然后,handler 会为其绑定上运行中的监听器做好替换准备(prepare_swap);如果该绑定上没有监听器,就绑定一个新的(listener::bind)并与之配对(Pending::pair);check跳过最后这一步。计划保留的入站,只要其监听器正在运行,且准入与运行中的不同,仍会得到一个新的用户表(先user_table,再prepare_users)。准入不同指的是:用户不同、principal 被替换(Arc::ptr_eq),或凭据不同(same_admissions)。
准备阶段返回一个
Prepared值(见下文)。某一步失败时,准备阶段返回该错误,已构建的内容随之被 drop,包括它绑定的监听器(以及 Unix 监听器的 socket 文件)。 - DNS:除非计划中有
-
提交(
Actor::commit)不会失败。它先采用用户 key(commit_keys),发布限速,并替换目标。当计划中有Step::PublishPlane时,它发布 plane(epoch加一),取消计划未复用的负载均衡器的探测 token,并在根 token 的一个子 token 下启动新负载均衡器的探测,这些探测用新解析器的服务器(dns.servers)解析,用TcpDialer::new(self.socket.clone())拨号;没有发布时,不动探测 token。无论是否发布,它都保存路由表、解析器和 DNS 服务目标。然后它撤销被移除的 principal,启动新的监听器(每个以root.child_token()作为停止 token),替换 handler,存入用户表,并遍历各个步骤:Step::StopAccepting停止该绑定上的监听器,Step::CloseSessions关闭该入站的会话,Step::Drain把旧版本交给它的排空策略。最后它保存准入和运行状态,并以 debug 级别记录applied: {} built, {} reused, {} swapped, {} drained。
每一步的完整顺序见规划并应用变更。
准备阶段构建的内容
Section titled “准备阶段构建的内容”struct Prepared<U: UserId> { dns: Dns, dns_target: Option<Arc<Target>>, targets: HashMap<CompactString, Arc<Target>>, new_balancers: Vec<(CompactString, Arc<Balancer>, Duration, Duration)>, routes: Arc<CompiledRoutes>, admissions: HashMap<CompactString, Admissions<U>>, keys: UserKeys<U>, revoked: Vec<(Vec<Arc<Principal>>, UserRemovalPolicy)>, started: Vec<(BindSpec, Pending)>, swaps: Vec<(usize, Swap)>, users: Vec<(usize, Users)>,}| 字段 | 内容 |
|---|---|
dns、dns_target |
解析器和 DNS 服务目标,重建的或沿用的 |
targets |
spec 中的每个出站和负载均衡器 tag,按 tag 索引 |
new_balancers |
每个新构建的负载均衡器及其 tag、probe_interval 和 probe_timeout,提交阶段会启动它们的探测 |
routes |
路由表,新编译的或沿用的 |
admissions、keys |
每个入站准入了谁,以及暂存的用户 key |
revoked |
按入站列出不再被准入的 principal,以及适用于它们的移除策略 |
started |
新的绑定,每个都已与其 handler 配对(Pending) |
swaps |
为运行中监听器准备的 handler,按监听器下标索引 |
users |
为运行中监听器准备的用户表,按监听器下标索引 |
没有经过提交就被 drop 的 Prepared(check 就是用 map(drop) 这样做的)以同样的方式释放它持有的资源。
Actor::publish_speed_limits 在每次提交以及每次用户编辑的提交开始时运行,从不在准备阶段运行,所以被拒绝的变更不会影响现行的限速。它遍历 spec 中每个用户集的每个用户,按 UserKey 收集每个用户的 speed_limit。同时属于多个用户集的用户取其中最小的限速;没有 key 的用户(没有任何入站准入过的用户)被跳过。随后 Tracker::set_speed_limits 为列表中每个用户的 pacer 设置速率,并把它持有的其他所有 pacer 恢复为不限速,所以存活的流会在下一次 poll 时感受到变化。pacer 见流跟踪、统计与限速。
set_users、upsert_user 和 remove_user 不做规划。Actor::edit_users 在运行中 spec 的副本里修改一个用户集,只重建准入该用户集的那些入站的用户表:
-
取出运行中的 spec(没有时,编辑以
ApplyError::Stopped失败),并找到该用户集。用户集不存在时为针对Resource::UserSet的ApplyError::Invalid,显示为user set {tag}: no such user set。 -
执行编辑。移除不在该用户集中的用户返回
Ok(false),不做任何改动。Set和Upsert总是算作一次变更。 -
对
users指名该用户集的每个入站执行准备:- 用
validate_admission针对该入站检查这个用户集,这是用户编辑唯一运行的校验规则; - 在 key 的暂存副本上准入其用户;
- 按绑定找到该入站的监听器,使用
expect("every running inbound has its listener"); - 构建用户表(
user_table),并检查它是否适用于监听器当前 handler 所服务的协议(Listener::prepare_users→UserTable::fits)。不匹配时,编辑以ApplyError::Build失败,错误为the handler is not of the kind its listener serves(io::ErrorKind::InvalidInput)。
第一个失败就会拒绝整个编辑,不存入任何表。
- 用
-
提交:采用 key,发布限速,按各入站的移除策略撤销已不存在的 principal,然后存入每张表及其准入,并把编辑后的 spec 保存为运行中的 spec。
用户编辑从不发布 plane,所以 epoch 不变。与应用一样,被移除用户的 principal 在存入任何新用户表之前就被撤销;第一个流出示已撤销 principal 的会话,在任何策略下都会被关闭(Session::admit)。细节见用户、principal 与会话。
token、任务与任务树
Section titled “token、任务与任务树”取消 token
Section titled “取消 token”flowchart TB root["根 token(Actor::new)"] stop["监听器 stop token,每个监听器一个"] tun["TUN 运行时 token,每个运行时一个"] probe["探测 token,每个新负载均衡器一个"] session["会话 token,每个会话一个"] bg["background_stop(独立)"] closed["Target closed token,每个出站版本一个(独立)"] root --> stop stop --> tun root --> probe root --> session
| token | 创建者 | 父 token | 取消者 | 停止的内容 |
|---|---|---|---|---|
root |
Actor::new |
无 | shutdown,在宽限期结束之后 |
它下面的一切;排空宽限定时器;仍在排空的 Hysteria 2 endpoint |
监听器 stop |
commit:为每个启动的监听器创建 root.child_token() |
root |
Step::StopAccepting、shutdown 或根 token |
accept 循环;Sessions::open 拒绝为它打开新会话;它下面的 TUN 运行时 |
| TUN 运行时 token | run_tun_inbound:每个运行时一个 stop.child_token() |
其监听器的 stop |
Listener::swap 发送的重启、其监听器的 stop,或其监听器句柄被 drop(重启通道关闭) |
一个 TUN 运行时及其服务的每个流 |
| 探测 token | commit:为每个新构建的负载均衡器创建 root.child_token() |
root |
丢弃或重建该负载均衡器的发布、shutdown 或根 token |
该负载均衡器的探测任务 |
| 会话 token | Sessions::open:root.child_token() |
root |
close、移除策略、Step::CloseSessions、Session::admit 拒绝已撤销的 principal、会话自身结束(Session::drop),或根 token |
在流式监听器上,是服务该连接的任务;Hysteria 2 连接(由协议关闭);等待它的移除宽限定时器 |
background_stop |
Actor::new:CancellationToken::new() |
无 | shutdown,在每个被跟踪的任务结束之后 |
采样器和用量推送任务 |
Target closed token |
Target::outbound 和 Target::balancer:CancellationToken::new() |
无 | Target::close_flows,在排空策略为 Close 或 CloseAfter 时 |
在该出站版本上打开的每个 stream 和数据报 link(Guarded)。经负载均衡器发送的流,由承载它的成员的 token 守护。 |
会话 token 挂在根 token 之下,而不是挂在接受该连接的监听器之下。正因如此,一次应用可以停止监听器、重新绑定它或为它更换 handler,而不触及它的连接:停止监听器只是停止接受连接,别的什么都不做。真正会结束会话的情况列在一次应用保留什么中。
最后一个 Arc<Session> 被 drop 时,会话结束。Session::drop 取消它的 token,从而结束仍在等待它的移除宽限定时器;在其用户的账本账户中关闭该会话,并入尚未统计的字节(Account::close);然后把它从注册表中移除。会话 id 从 1 开始递增(Sessions::next,AtomicU64::new(1))。
background_stop 有意不作为根 token 的子 token。采样器和用量 sink 必须比每个连接都活得久,才能统计到连接的最后一批字节,而取消根 token 正是结束连接的手段。
supervisor 在三个地方运行它的任务:普通的 tokio::spawn、它唯一的 TaskTracker,以及阻塞线程池。
| 任务 | spawn 方式 | 由谁 spawn | 何时结束 |
|---|---|---|---|
| actor | tokio::spawn |
SupervisorBuilder::start |
Shutdown 处理完毕,或所有句柄都已不存在 |
| 采样器 | tokio::spawn,保存在 background 中 |
SupervisorBuilder::start |
background_stop,在最后一个 tick 之后 |
| 用量推送 | tokio::spawn,保存在 background 中 |
SupervisorBuilder::start,设置了 sink 时 |
background_stop,在最后一次调和与取出之后,非空时上报 |
| 流式 accept 循环 | TaskTracker::spawn |
Listener::start |
其监听器的 stop。select! 以 biased 优先处理 stop;当 Sessions::open 发现 stop 已取消时,循环也会退出。 |
| Hysteria 2 监听器 | TaskTracker::spawn |
Listener::start |
其 stop 被取消且最后一个连接已结束,或根 token 被取消 |
| TUN 运行时循环 | TaskTracker::spawn |
Listener::start,以及运行时已经结束的替换 |
其监听器的 stop、监听器句柄被 drop,或其运行时自行结束 |
| 一个被接受的 socket | TaskTracker::spawn,与该 socket 的会话 token 一起放在 select! 中(Shared::spawn_until) |
流式 accept 循环 | 其传输层结束,或其会话 token 被取消。它运行传输层;Unix socket 的连接就在这个任务上服务。 |
| TCP 传输层产生的一条字节流 | 在其会话 token 下的 Shared::spawn_until |
socket 的任务,从传输层的 accept 回调中 | 该流的运行时结束,或其会话 token 被取消。这个会话是 socket 自己的会话;在 gRPC 承载连接上,则是为这条流打开的会话。 |
| 一个 Hysteria 2 连接 | 其监听器任务内的一个 JoinSet(Hy2Inbound::run) |
Hysteria 2 监听器 | 连接结束,或在其会话 token 被取消时由协议关闭 |
| 负载均衡器探测,每个成员一个 | TaskTracker::spawn |
Balancer::spawn_probe,由 commit 调用 |
它的探测 token |
| 移除宽限定时器 | TaskTracker::spawn |
CloseAfter 下的 Sessions::revoke |
宽限期到期,或会话 token 被取消 |
| 排空宽限定时器 | TaskTracker::spawn |
CloseAfter 下的 drain |
宽限期到期,或根 token 被取消 |
| 一个采样器 tick、一次账本调和 | spawn_blocking |
采样器、用量推送 | tick 或调和完成 |
只有一个 TaskTracker,它在 Actor::new 中创建,由 Shared 和会话注册表共享,所以 shutdown 用一次 wait 就能等待每个 accept 循环、连接、探测和定时器;Hysteria 2 监听器的任务则等待它自己的连接。Shared::spawn_until 这样包装每个 socket 和每条流:
self.tracker.spawn(async move { tokio::select! { _ = token.cancelled() => {} _ = fut => {} }});所以取消会话 token 会在下一次 poll 时 drop 这个 future,连同它的 socket。在 gRPC 承载连接(carrier)上,每个 HTTP/2 stream 都是一个独立的会话。在其监听器的 stop 被取消之后才完成 QUIC 握手的 Hysteria 2 连接得不到会话:run_hysteria_inbound 交给它一个已经取消的 token,所以没有会话能逃过被移除入站的关闭。accept 循环及其服务的内容见监听器与服务循环。
socket 策略
Section titled “socket 策略”pub fn socket_options(self, socket: SocketOptions) -> Self;SocketOptions(environment/src/dial/socket.rs)是作用于拨号器打开的每个 socket 的策略:源地址 bind_address(只作用于同一地址族的 socket)、interface、数据包 mark(SO_MARK)、tcp_keepalive 空闲时间,以及在每个 socket 创建之后、绑定或连接之前运行的 hook。各字段及平台支持情况见拨号器与 socket 策略。
builder 把策略存进 actor,actor 再把它传给每个会打开出站 socket 的构建器:
| socket | 通过什么获得策略 |
|---|---|
| 代理出站的每次拨号,包括 TCP 及其传输层 | build/outbound.rs → Dialer::new(socket) |
| 直连(freedom)流的 TCP 拨号和 UDP socket | build/outbound.rs → Dialer::new(socket),其 TCP 和 UDP 两半都带着该策略 |
| SOCKS 出站的 UDP 中继 socket | build/outbound.rs → UdpDialer::new(socket) |
| Hysteria 2 出站的 QUIC socket | Hy2Config::socket |
| WireGuard 出站的隧道 socket | WgConfig::socket |
| 解析器从本机发出的每个查询 | build/dns.rs → ResolverOptions::socket |
| 负载均衡器健康探测 | commit → TcpDialer::new(self.socket.clone()) |
策略在 supervisor 的整个生命周期内有效;应用无法改变它。它不作用于:
- 监听器 socket,它们由
system::listener::bind用标准库打开:策略针对的是本进程发起的流量; - 系统解析器。
getaddrinfo在 C 库内部打开它的 socket,任何策略都够不着。后端为系统解析器的解析器(例如DnsSpec::Single(Backend::System))使用的就是它。
移动端 VPN 正是在这里把代理自身的流量排除在隧道之外。etemenanki-ffi 传入一个 hook,它调用 app 的 Platform::protect(fd)(在 Android 上是 VpnService.protect);平台拒绝的 socket 以 the platform refused to protect the socket 失败,它所服务的那次拨号也随之失败。任何类型的 hook 错误都会让拨号失败:不会有数据从策略未接受的 socket 发出。
check 使用 SocketOptions::default() 构建:它从不拨号,所以不需要策略。
pub async fn shutdown(&self, grace: Duration);shutdown 发送 Command::Shutdown(grace, reply),actor 运行 Actor::shutdown:
- 取消每个监听器的
stoptoken。accept 循环返回并释放它们的 socket(Unix 监听器会删除它的 socket 文件);Sessions::open拒绝为它们打开任何会话。Hysteria 2 监听器停止接纳新连接并开始排空。TUN 运行时的 token 是其监听器stop的子 token,所以每个 TUN 运行时及其服务的流在这里就结束了,早于宽限期:TUN 流不是会话。 - 取消每个负载均衡器的探测 token。
- 关闭
TaskTracker,使wait能在它清空后结束。 - 最多等待
grace,让每个被跟踪的任务自行结束:tokio::time::timeout(grace, tracker.wait())。 - 取消根 token。每个会话 token 都是它的子 token,所以仍在服务流式连接的每个任务会在下一次 poll 时结束,协议会关闭每个仍然打开的 Hysteria 2 连接;移除和排空宽限定时器结束;Hysteria 2 监听器在排空过程中关闭它的 endpoint,所以仍处于 QUIC 握手中的客户端(它还没有会话)也会结束。
- 等待每个被跟踪的任务结束。
- drop 每个
Listener句柄。 - 取消
background_stop,并等待采样器和用量推送任务。此时每个连接都已结束并并入了它最后的字节,所以采样器运行最后一个 tick,推送任务再调和并取出一次,非空时上报这一批。
随后 actor 应答并返回,drop 命令接收端。grace 是存活连接可用于结束的时长;Duration::ZERO 会立即关闭它们。
所有句柄都在没有关停的情况下被 drop 时,循环以 Duration::ZERO 运行同样的流程。通过 tracker() 取得的 Tracker 克隆不会让 actor 保持存活。
关停之后:
| 调用 | 结果 |
|---|---|
apply、apply_with、update、update_with |
Err(ApplyError::Stopped):the supervisor has shut down |
set_users、upsert_user、remove_user |
Err(ApplyError::Stopped) |
close |
0 |
sessions |
空 Vec |
shutdown |
立即返回 |
epoch |
最后发布的 epoch |
tracker |
仍可读取;最后的 StatsSnapshot 来自最后一个 tick |
take_usage、usage_snapshot |
仍然可用:没有被 sink 取走的字节仍可取出 |
排在 Shutdown 之后的命令同样得到 Stopped:它的应答发送端随通道一起被 drop。
supervisor/src/entity/id.rs 按值而不是按指针为正在运行的东西命名。fan-out UDP link 用打开每个子 link 时所在的 OutboundId 作为该子 link 的 key,而不是用 Arc 的地址,因为出站的寿命长于单次应用,同一个 tag 可能有一个旧版本在新版本旁边排空,而已释放的内存分配可能被其后继复用。
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]pub struct OutboundId { pub tag: CompactString, pub version: u64,}
pub enum ResourceKind { Inbound, Outbound, Balancer, UserSet }
pub enum Resource { Inbound(CompactString), Outbound(OutboundId), Balancer(CompactString), UserSet(CompactString), Route, Dns,}
pub struct SessionId(u64);pub struct UserKey(u64);pub struct FlowId(u64);pub struct RuleId(u32);| id | 指称 | 分配者 | 首个值 | Display |
|---|---|---|---|---|
OutboundId |
一个出站或负载均衡器 tag 的一次构建 | plan:对还没有版本的 tag 为 1,每次重建取上一版本加一 |
1 |
{tag}@v{version} |
ResourceKind |
一种带 tag 的资源,用于 DuplicateTag 和 UnknownReference |
— | — | inbound、outbound、balancer、user set |
Resource |
一个资源,用于 Plan、ApplyReport 和 ApplyError |
— | — | inbound {tag}、outbound {tag}@v{n}、balancer {tag}、user set {tag}、route、dns |
SessionId |
一个被接受的连接:一个 socket、一个 QUIC 连接,或 gRPC 承载连接的一条 stream | Sessions::open,一个 AtomicU64 计数器 |
1 |
session#{n} |
UserKey |
一个被某个入站准入的用户 id | UserKeys::key,由 admit 在某个入站第一次准入该 id 时调用 |
1 |
无 |
FlowId |
一个已选路的流:一条 stream,或 UDP 关联的一个子 link | Tracker,一个 AtomicU64 计数器 |
1 |
flow#{n} |
RuleId |
为流选路的那个 plane 中,某条规则在 RouteSpec::rules 里的下标 |
CompiledRoutes,每条规则一个 |
0 |
无 |
- 版本从不回退。
RunningState::versions保留已移除 tag 的版本,所以重新加回的 tag 得到的版本是它的旧构建(可能仍在排空)从未用过的。负载均衡器也在同一个 map 中编号,负载均衡器目标同样带有OutboundId,尽管Resource::Balancer只包含 tag。 - 空 tag 归 supervisor 自己使用。 校验会拒绝任何类型的空 tag(
a {kind} tag must not be empty),所以 supervisor 把它用于自己创建的目标(supervisor.rs中的internal_id(version)构建一个 tag 为空的OutboundId):空 plane 的 blackhole(版本0,在第一次应用之前)和 DNS 服务目标。后者只在 DNS 被重建时于准备阶段构建,使用即将承载它的 plane 的 epoch(internal_id(epoch + 1));否则保留运行中的那个。采样器把 DNS 服务报告在空的出站 tag 之下。 - 会话 id 和流 id 在 supervisor 的生命周期内唯一,从不复用。
SessionId、UserKey、FlowId和RuleId以const fn形式提供new(raw)和get(),所以前端程序可以把 id 交出去再收回来,就像 FFI 在close_connection中用FlowId::new(id)所做的那样。 - 用户被移除后,其 key 仍保持分配,所以
Keep策略留下的会话仍能通过Selector::User找到,它们最后的用量仍以该用户的 id 上报。只有当某个入站准入该用户时,也就是用户具有该入站协议所接受的凭据类型时,才会分配 key;没有任何入站准入的用户没有 key。keyn指称第n个被分配的 id(id_of:ids[n - 1]);UserKey(0)不指称任何人。 - 规则 id 属于某一个 plane。 对走默认路由的流,以及被 DNS 服务拦截的流,流的规则(
Routed::rule、FlowEntry::rule)为None。路由器用expect("fewer than 2^32 rules")把下标转换为u32。
前端程序如何使用它
Section titled “前端程序如何使用它”| etemenanki-app | etemenanki-ffi | katana | |
|---|---|---|---|
| 代码 | app/src/instance.rs → Core、Instance |
ffi/src/proxy.rs → Proxy,运行 app 的 Core |
src/manager/node.rs → NodeManager |
U |
UserName(app/src/lower.rs → UserKey) |
UserName,通过 app 的 lowering |
Uid(i64),面板的用户 id;name() 是该数字的 UserName::Username |
| supervisor 数量 | 每个进程一个 | 每个启动的代理一个 | 每个节点一个 |
| 启动 | Supervisor::builder(),全部取默认值 |
Supervisor::builder().socket_options(...),带 protect hook |
Supervisor::start(spec),全部取默认值 |
| 变更 | 每次重载时 apply(spec) |
通过 app 的 Core::reload_with 调用 apply(spec) |
apply_with(spec, ApplyOptions { allow_disruptive: true }),spec 由面板的应答降为 |
| 读取 | tracker(),供 REST API 使用 |
tracker():stats()、flows()、kill(FlowId) |
take_usage() 用于流量上报,tracker().subscribe() 用于审计 |
| 关停 | 收到信号时 shutdown(SHUTDOWN_GRACE),5 s |
在 Proxy::stop 中 shutdown(STOP_GRACE),2 s;然后给运行时的线程 RUNTIME_GRACE,3 s |
节点停止时 shutdown(Duration::ZERO),然后最后一次 take_usage,用于最终上报 |
- etemenanki-app 把 TOML 配置及其订阅文件降为 spec(
lower),在--test背后以及 REST API 的路由切换写入编辑后的配置之前运行check,并以默认选项应用每次重载。它在第一次应用后记录config loaded: {summary},在每次成功应用的重载后记录config reloaded: {summary},其中 summary 列出ApplyReport中每个非空的分组(built [...]; swapped [...]等),或者为nothing to run。中断性变更会被拒绝,并记录为reload refused, keeping the running config: inbound {inbound}: {reason}, which would end its live connections; restart to apply it。它的 lowering 在etemenanki-app:从 TOML 到 spec中介绍,它的重载在etemenanki-app:运行、重载与关停中介绍。配置的 REST API 在启动时无法绑定时,二进制会以Duration::ZERO关停 supervisor 并退出。 - etemenanki-ffi 用它自己的 lowering(
ffi/src/client.rs→lowering)运行 app 的Core。它拒绝向他人提供代理协议服务的入站:所有服务端协议,以及不在回环地址或 Unix socket 上的 SOCKS 或 HTTP 入站(ServerInbound)。它拒绝多于一个的 TUN 入站,并把 TUN 入站以BindSpec::Tun(TunSource::Fd(...))的形式绑定到手机提供的描述符上;配置中没有 TUN 入站时,它会添加一个 tag 为tun(TUN_TAG)、使用协议默认值(DEFAULT_MTU、DEFAULT_UDP_IDLE_TIMEOUT、DEFAULT_MAX_FLOWS,开启嗅探)的入站。它不提供配置中的[api]。它的 socket 策略就是本页所述的那个。它在移动端库(etemenanki-ffi)中介绍。 - katana 每个节点运行一个 supervisor,并在设置
allow_disruptive的情况下,应用由面板的节点信息和用户列表降为的 spec 。当节点的传输层或协议设置与上次应用的不同时(NodeInfo::transport_eq、NodeInfo::protocol_eq,它们涵盖端口、TLS、混淆、节点类型、加密方法和服务端密钥等),或者当面板应答无法体现的本地配置修改要求这样做时,它随后调用close(Selector::All)。它的审计任务用tracker().subscribe()订阅流事件,并从FlowEvent::Opened记录审计命中。见节点管理器和流量计费。
这个 crate 遵循几条规则,本部分的其他页面都依赖它们。总体原则见设计原则。
- 每条规则只有一个所有者。 关于 spec 中各个带类型部分如何组合的每一条规则都位于
build/validate.rs。前端程序负责只有它自己才知道的东西:解析自己的字符串、读取自己的文件、自己 schema 的默认值、如何为用户命名,以及 spec 无法表达的键组合。正如app/src/lower.rs所说,一条规则如果在前端程序里也检查一遍,就会成为一份会逐渐走样的副本;而只在前端程序里检查的规则,对其他每个前端程序来说都是缺失的。 - 用
==比较带类型的 spec。Spec保存的是带类型的值(Destination、Uuid、解码后的密钥字节、PEM 字节),而不是原始配置文本。这让规划器可以凭相等性决定保留什么、重建什么,而无需知道两个 spec 各自从哪里来。 - 先准备,再提交。 所有可能失败的工作都发生在任何正在运行的东西改变之前,所以被拒绝的 spec 让运行状态保持原样。
- 读-复制-更新(RCU)式路由。 路由和出站通过一次原子 store 一起发布,所以规则永远不会指名 plane 中没有的 tag;旧 plane 在最后一个持有其某部分的流结束时被释放。
- 用身份标识代替指针。 出站版本、会话、用户、流和规则都按值命名(见实体 id)。
- 慢的观察者不会拖住数据平面。 每个流的计数器归它自己;采样器每个 tick 在阻塞线程池上汇总一次;流事件通过一个有界的广播通道(
EVENT_CAPACITY,1024)发出,落后的订阅者会丢失事件,而不是拖住发送方;用量 sink 在自己的任务上运行。没有任何按字节的操作经过 actor。
| 不变量 | 由什么保证 | 由哪个测试固定 |
|---|---|---|
| 被拒绝的 spec 不改变任何东西:监听器、plane 和用户保持原样 | 准备阶段不改动任何正在运行的东西;提交不会失败;被拒绝的准备阶段绑定的监听器在它返回错误时被释放 | a_refused_spec_changes_nothing(一次校验拒绝和一次绑定失败) |
| 应用和用户编辑从不并发执行 | 一个 actor 任务,一次一个命令 | 没有直接的测试 |
| 保留其入站的应用不会结束连接 | 会话 token 是根 token 的子 token;未变的绑定保留其监听器 | an_established_connection_survives_an_apply_that_keeps_its_inbound |
| 路由变更会作用于下一个流、mux 子流和 UDP 包,但不会为已打开的 TCP 流重新选路 | connector 每次 connect 加载一次 plane;fan-out 按包选路 |
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 |
epoch 统计已发布的 plane;被拒绝的应用不发布任何 plane |
只在提交阶段的 Step::PublishPlane 处 epoch += 1 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow(epoch + 1)、a_refused_spec_changes_nothing(不变) |
| 第一次应用之前的 plane 丢弃一切 | Plane::empty,epoch 为 0 |
the_empty_plane_drops_everything |
| spec 和解析器都未变的出站被沿用;重建的出站是一个新版本 | Step::Reuse 放入运行中的 Arc<Target>;Step::Build 取下一个版本;DnsSpec 变化会重建每个非 blackhole 出站 |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply、each_rebuild_of_a_tag_takes_the_next_version、a_tag_added_back_takes_a_version_it_never_had |
Step::Disrupt 变更需要 allow_disruptive |
在准备之前检查 Plan::disruptions |
a_hysteria2_obfuscation_change_needs_allow_disruptive、a_tun_device_change_needs_allow_disruptive |
| 同一绑定上的协议变更保留 socket | 监听器以 BindSpec 为 key;替换 handler |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
| 每个会话都是根 token 的子 token | Sessions::open → root.child_token() |
cancelling_the_root_cancels_every_session |
| 已停止的监听器不打开会话 | Sessions::open 在注册表的锁下检查 stop |
a_stopped_listener_opens_no_session |
| 移除入站会关闭它的会话 | 先 Step::StopAccepting,再 Step::CloseSessions |
removing_an_inbound_closes_its_sessions |
close 只统计它关闭的会话 |
Sessions::close 跳过已取消的 token |
a_session_already_closed_is_not_counted_again |
| 关停不会等满宽限定时器 | 定时器运行在 TaskTracker 上;移除定时器在会话 token 上 select,排空定时器在根 token 上 select,而关停会取消根 token |
shutdown_ends_pending_grace_timers |
| 关停不会等待停滞的 Hysteria 2 握手 | 根 token 一被取消,run_hysteria_inbound 就关闭 endpoint |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
| 最后一批用量在最后一个会话结束之后取出 | background_stop 在 tracker.wait() 之后被取消 |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
| 慢的 sink 永远不会拖住连接 | sink 在自己的任务上运行 | a_slow_usage_sink_does_not_stall_the_data_plane |
| 每个出站 socket 都在 socket 策略下打开 | actor 把 SocketOptions 传给每个构建器和探测 |
direct_flows_are_dialed_on_hooked_sockets、a_socks_upstream_is_dialed_on_hooked_sockets、a_refusing_hook_fails_the_dial |
| 用户 key 在任何会话能够绑定它之前发布 | commit_keys 在提交和用户编辑的提交中最先运行 |
没有测试;否则 UsageBook::name 会 panic |
| 限速只在应用或用户编辑提交时改变 | publish_speed_limits 只在提交中调用 |
a_speed_limit_change_keeps_the_users_sessions(会话被保留) |
| spec 中没有空 tag,所以内部目标永远不会与 spec 的目标冲突 | validate 中的 unique_tags |
no_tag_may_be_empty |
失败路径与取消
Section titled “失败路径与取消”start 只会因第一次应用而失败,错误与之后的应用可能返回的相同,只是没有 Disruptive(此时还没有东西在运行)和 Stopped。失败时还没有 spawn 任何东西:采样器、用量推送和 actor 都只在第一次应用提交之后才 spawn。被拒绝的准备阶段中绑定的监听器,在准备阶段返回错误时被释放。start 会断言 sample_interval 不为零(the sample interval must not be zero)。
被拒绝的变更
Section titled “被拒绝的变更”被拒绝的应用或用户编辑返回它的 ApplyError,运行中的 spec、版本、plane、监听器、用户表、key 和限速都保持原样。准备期间分配的用户 key 位于一个暂存副本(self.keys.clone())上,只在提交时才被采用,所以被拒绝的变更不会分配任何 key。
不再等待的调用方
Section titled “不再等待的调用方”命令一旦发出,drop 调用的 future 并不会取消它:actor 会完成这个命令,并丢弃应答(let _ = reply.send(...))。调用方放弃等待的应用仍会提交。在等待通道空位时被 drop 的 future 则什么也不发送。
从外部结束的连接
Section titled “从外部结束的连接”close、移除策略、Step::CloseSessions 和关停都通过会话 token 结束连接。在流式监听器上,服务该连接的每个任务都通过 Shared::spawn_until 运行在这个 token 之下,所以取消它会在下一次 poll 时 drop 任务的 future。Hysteria 2 连接在其会话 token 被取消时由协议关闭(见监听器与服务循环)。
Tracker::kill 和 Tracker::kill_where 按流强制关闭(kill):它们设置流的 kill 标志并唤醒它,所以它在任一方向的下一次 poll 都会失败;它所属的会话不受影响。排空策略在比会话低一层的地方工作:它取消出站版本的 closed token,于是在该版本上打开的每个 Guarded stream 或数据报 link,下一次读或写都会以 the outbound this flow was opened on was drained(ConnectionAborted)失败,中继它的运行时随之结束这个流。路由一侧见 plane:为每个流选路。
每个经过 actor 的调用都返回 Stopped、0、空列表或什么也不返回,如关停中所列。绕过 actor 的读取继续应答。
| 限制 | 值 | 位置 | 含义 |
|---|---|---|---|
| 命令通道 | 16 个命令 | SupervisorBuilder::start → mpsc::channel(16) |
已有 16 个命令排队时,调用方的 send 等待 |
DEFAULT_SAMPLE_INTERVAL |
1 s | track/sampler.rs |
采样器的 tick 间隔,因而也是 take_usage 对存活会话的滞后 |
EVENT_CAPACITY |
1024 个事件 | track/mod.rs |
订阅者在开始丢失事件之前最多可以落后的流事件数 |
sample_interval |
不能为零 | SupervisorBuilder::start 中的断言 |
否则 panic |
| 出站版本 | u64,每个 tag 从 1 开始 |
plan |
已移除的 tag 也保留,所以 RunningState::versions 为每个出现过的 tag 保留一个条目 |
SessionId、FlowId |
u64,从 1 开始 |
Sessions::open、Tracker |
在 supervisor 的生命周期内唯一 |
UserKey |
u64,从 1 开始 |
UserKeys::key |
每个用户 id 一个,在某个入站第一次准入它时分配 |
RuleId |
u32 |
topology/router.rs |
每个路由表少于 2^32 条规则 |
| etemenanki-app 关停宽限期 | 5 s | app/src/main.rs → SHUTDOWN_GRACE |
|
| etemenanki-ffi 停止宽限期 | 2 s | ffi/src/proxy.rs → STOP_GRACE |
|
| etemenanki-ffi 运行时宽限期 | 3 s | ffi/src/proxy.rs → RUNTIME_GRACE |
在停止宽限期之后,留给运行时的线程(shutdown_timeout) |
| katana 关停宽限期 | 0 | src/manager/node.rs |
按入站的连接数限制和 accept 出错时的退避属于服务层,见监听器与服务循环。整个系统的限制汇总在限制、超时与内存中。
集成测试在回环 socket 上启动真实的 supervisor,并用 supervisor/tests/support/mod.rs 中的最小协议客户端驱动它们:
| 辅助函数 | 作用 |
|---|---|
start_on、start_built |
在为新选取的端口构建的 spec 上启动 supervisor;遇到 ApplyError::Bind 时重新选取,最多重试 5 次,因为端口可能在选取与绑定之间被占用。start_built 接受一个 builder 工厂,socket 策略测试和用量 sink 测试就是通过它安装各自的设置的。 |
base_spec |
一个 direct freedom 出站、一个 block blackhole、默认走 direct 的路由、DnsSpec::Single(Backend::System) 以及默认策略。 |
QUIET、SOON |
等待不应发生的事情用 400 ms,等待应该发生的事情用 5 s。 |
| 客户端 | SOCKS connect 与 UDP associate、HTTP CONNECT、一个 VLESS mux 客户端、回显服务器、一张自签名证书。 |
hot-swap 和 tracking 测试最后会对它们启动的每个 supervisor 调用 shutdown(Duration::ZERO)。socket 策略测试则直接 drop 句柄,这会以零宽限期运行同样的关停流程。
| 测试 | 文件 | 固定的行为 |
|---|---|---|
an_established_connection_survives_an_apply_that_keeps_its_inbound |
tests/hot_swap.rs |
被保留的入站列在 reused 中,而不是 swapped 或 rebound,它的连接继续传输数据 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
tests/hot_swap.rs |
epoch 加一;下一个 UDP 包和新的 TCP 流走新路由;已打开的 TCP 流不走 |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
tests/hot_swap.rs |
变更之后打开的 mux 子流走新路由;之前的子流保持自己原来的路由 |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close |
tests/hot_swap.rs |
在 freedom 出站 direct 上:drained 列出旧版本;在 Keep 下,在旧版本上打开的流继续回显,在 Close 下它被关闭;新流使用新版本 |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
tests/hot_swap.rs |
被复用的 Hysteria 2 出站不会再次握手;重建的出站重新拨号,而旧出站上的流继续回显 |
a_hysteria2_obfuscation_change_needs_allow_disruptive |
tests/hot_swap.rs |
没有标志时为 Disruptive;有标志时,restarted 列出该入站,端口保持不变 |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
tests/hot_swap.rs |
列在 swapped 中而不是 rebound;替换之前接受的连接由旧 handler 服务 |
a_refused_spec_changes_nothing |
tests/hot_swap.rs |
未知引用和绑定失败都不改变 epoch、用户和路由;被拒绝的准备阶段中绑定的监听器被释放 |
removing_a_user_closes_only_their_sessions_and_refuses_them_after |
tests/hot_swap.rs |
在 SOCKS 入站上,分别在 Close 和 Keep 下:被移除用户的会话按策略被关闭或保留,另一个用户的会话不受影响,被移除用户的新登录被拒绝;close(Selector::User) 返回 1 |
a_speed_limit_change_keeps_the_users_sessions |
tests/hot_swap.rs |
限速变更前后的会话 id 相同 |
removing_an_inbound_closes_its_sessions |
tests/hot_swap.rs |
removed 列出它;它的连接被关闭;它的端口拒绝连接;另一个入站不受影响 |
a_tun_device_change_needs_allow_disruptive |
tests/hot_swap.rs |
没有标志时,设备变更和设置变更都以 Disruptive 被拒绝。需要 CAP_NET_ADMIN;无法创建设备时跳过(并算作通过) |
shutdown_ends_pending_grace_timers |
tests/hot_swap.rs |
一小时的移除和排空定时器不会拖住 shutdown(Duration::ZERO) |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
tests/hot_swap.rs |
QUIC 握手永不完成的客户端不会拖住关停 |
direct_flows_are_dialed_on_hooked_sockets |
tests/socket_policy.rs |
直连的 TCP 和 UDP socket 都经过 hook |
a_socks_upstream_is_dialed_on_hooked_sockets |
tests/socket_policy.rs |
SOCKS 上游的控制连接及其 UDP 中继 socket 都经过 hook |
a_refusing_hook_fails_the_dial |
tests/socket_policy.rs |
hook 出错会让 SOCKS 请求失败 |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
tests/tracking.rs |
各批次大约相隔一个间隔到达;连同关停时的那一批,它们包含每一个字节;之后 take_usage 为空,usage_snapshot 仍然应答 |
a_slow_usage_sink_does_not_stall_the_data_plane |
tests/tracking.rs |
每次上报耗时 1 s 的 sink 只推迟它自己的批次 |
cancelling_the_root_cancels_every_session |
tests/unit/session.rs |
取消根 token 时,已绑定和未绑定的会话都会结束 |
a_stopped_listener_opens_no_session |
tests/unit/session.rs |
在已取消的 stop 下,Sessions::open 不登记任何会话 |
the_empty_plane_drops_everything |
tests/unit/plane.rs |
epoch 为 0;TCP 和 UDP(包括 53 端口)都进入 blackhole |
the_first_plan_binds_and_builds_everything_then_publishes |
tests/unit/plan.rs |
第一次计划没有中断性步骤,并以 PublishPlane 结束;每个 tag 从版本 1 开始 |
each_rebuild_of_a_tag_takes_the_next_version、a_tag_added_back_takes_a_version_it_never_had |
tests/unit/plan.rs |
版本编号 |
no_tag_may_be_empty |
tests/unit/validate.rs |
空 tag 始终归 supervisor 使用 |
没有测试断言最后一个句柄被 drop 时的行为(socket 策略测试以这种方式结束,但不检查它),也没有测试断言 ApplyError::Stopped 或命令通道的容量。check 通过 app 得到覆盖,例如:a_member_with_no_upstream_is_refused(app/tests/integration/e2e_balancer.rs)对一个以 freedom 出站为成员的负载均衡器运行二进制的 --test,并预期失败;a_backend_without_its_server_is_rejected(app/tests/integration/e2e_dns.rs)对一个没有服务器的 UDP DNS 后端运行 --test,并预期失败;every_lowered_node_builds_and_the_dns_setup_with_it(app/tests/unit/subscribe.rs)对降为 spec 后的示例订阅文件调用 check。
supervisor/benches/tracking.rs 在一个服务器规模的账本上测量一次用量取出及为其提供数据的调和的开销,以及计量流的每字节开销(cargo bench -p etemenanki-supervisor --bench tracking);见流跟踪、统计与限速。