规划并应用变更
源码文件:41 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/topology/spec_plan/mod.rsEtemenanki/supervisor/src/topology/spec_plan/inbound.rsEtemenanki/supervisor/src/topology/spec_plan/outbound.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/policy.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/system/listener.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/build/users.rsEtemenanki/supervisor/src/build/inbound.rsEtemenanki/supervisor/src/build/outbound.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/build/route.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/lib.rsEtemenanki/protocols/src/dns/mod.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/support/mod.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/api.rsEtemenanki/ffi/src/proxy.rsEtemenanki/app/tests/integration/e2e_reload.rsEtemenanki/app/tests/integration/e2e_hysteria_inbound.rskatana/src/manager/node.rskatana/src/runtime.rs
supervisor(监管器)从不靠重启来改变它运行的内容。每一次变更,包括第一份 spec(期望状态),都经过同一个调和过程:plan 比较正在运行的内容与期望的 spec,列出两者之间的步骤;actor 先准备(prepare)所有可能失败的部分,但其中任何一样都不对外服务,然后提交(commit)其余不会失败的部分。在提交之前任何一处被拒绝的 spec,都会让监听器、数据平面、用户和限速保持原样。这取代了 etemenanki 2.x 的模型:那时每次重载都会构建一整套新的 generation(一代实例),包括监听器、出站和路由;调和则按资源进行,spec 没有变化的资源会继续运行,保留它的 socket 和状态。
本页面向修改 supervisor/src/topology/spec_plan/plan.rs、supervisor/src/supervisor.rs 中的应用(apply)路径,或 supervisor/src/build/apply.rs 的贡献者。内容包括:每类资源如何规划、出站版本与排空、中断性变更、准备和提交两个阶段的逐步过程、撤销用户、ApplyReport,以及 update。这个 crate 整体负责什么、actor 的其他命令见 supervisor 概览;spec 的类型见 spec:期望状态;plan 最先检查的规则见校验与 apply 错误。
| 阶段 | 代码 | 可能失败 | 可见效果 |
|---|---|---|---|
| 规划 | topology/spec_plan/plan.rs → plan |
是:validate 拒绝该 spec |
无。纯计算:不碰 socket、文件和时钟。 |
| 中断检查 | supervisor.rs → Actor::apply |
是:ApplyError::Disruptive |
无 |
| 准备 | supervisor.rs → Actor::prepare |
是:ApplyError::Build、ApplyError::Bind |
尚未对外服务任何东西。这里打开的 socket 和设备在提交之前不对外服务,被拒绝时会被关闭。 |
| 提交 | supervisor.rs → Actor::commit |
否 | 发布 plane(数据平面),启动并切换监听器,存入用户表,停止已不存在的资源,排空旧出站,撤销被移除的用户。 |
调和负责:
- 按资源决定是保留(
Reuse)还是构造(Build);对入站,还要决定是需要新的 socket(Bind),还是在旧 socket 上用SwapHandler换上新的 handler(处理器),还是什么都不用做; - 为出站或负载均衡器 tag 的每一次构建分配各自的版本,让旧构建可以在新构建旁边排空;
- 拒绝会结束存活连接的变更,除非调用方允许;
- 把路由和出站作为一个 plane 一起发布,且仅在 DNS、某个出站、某个负载均衡器或路由有变化时才发布;
- 把被替换或被移除的出站版本上的存活流交给它的排空策略处理,把被移除用户的存活会话交给入站的移除策略处理;
- 在
ApplyReport中报告它做了什么。
它不负责:
- 解析任何格式、读取配置文件或决定何时重载。前端程序把自己的格式降为
Spec,然后调用apply; - 定义什么是合法的 spec。
plan调用build/validate.rs→validate,并原样返回它的错误; - 服务连接或路由流。它只构建并切换做这些事的对象;见监听器与服务循环和 plane:逐流路由。
RunningState 与 Plan
Section titled “RunningState 与 Plan”#[derive(Debug, Clone)]pub struct RunningState<U: UserId> { pub spec: Option<Spec<U>>, /// Kept for tags that were removed too. pub versions: BTreeMap<CompactString, u64>,}
impl<U: UserId> RunningState<U> { pub fn empty() -> Self;}
#[derive(Debug, Clone, Default, PartialEq, Eq)]pub struct Plan { pub steps: Vec<Step>, /// The latest version of every outbound and balancer tag after the /// plan, removed tags included. pub versions: BTreeMap<CompactString, u64>,}
impl Plan { pub fn disruptions(&self) -> impl Iterator<Item = (&CompactString, &CompactString)>;}
pub fn plan<U: UserId>(running: &RunningState<U>, desired: &Spec<U>) -> Result<Plan, ApplyError>;RunningState是规划器看到的内容:最近一次应用的 spec(第一次应用之前为None),以及每个曾经应用过的出站和负载均衡器 tag 的最新版本。RunningState::empty()即{ spec: None, versions: empty }。versions从不裁剪。被移除的 tag 保留它的条目,因此重新加回的 tag 得到的版本号,是它那个可能仍在排空的旧构建从未用过的。Plan::disruptions()按步骤顺序,为每个Step::Disrupt产出(inbound, reason)。plan和这两个类型一样是公开的,但在此版本中没有前端程序直接调用它;它们都经由Supervisor::apply。
#[derive(Debug, Clone, PartialEq, Eq)]pub enum Step { Bind { inbound: CompactString, bind: BindSpec }, Build(Resource), Reuse(Resource), PublishPlane, SwapHandler { inbound: CompactString }, Disrupt { inbound: CompactString, reason: CompactString }, StopAccepting { inbound: CompactString, bind: BindSpec }, CloseSessions { inbound: CompactString }, Drain { outbound: OutboundId, policy: DrainPolicy },}| 步骤 | 含义 | 准备阶段的工作 | 提交阶段的工作 | 在报告中记为 |
|---|---|---|---|---|
Bind |
打开入站所要求的监听器:新入站,或 BindSpec 有变化的入站 |
listener::bind,然后 Pending::pair |
Listener::start |
rebound(该 tag 原本存在且绑定不同时) |
Build(r) |
按 spec 构造 r:新增的,或有变化的 |
构建 DNS 解析器、出站、负载均衡器、路由表或入站 handler | 发布它或把它换入 | built |
Reuse(r) |
保留正在运行的 r:它的 spec 没有变化 |
克隆它的 Arc(出站、负载均衡器、路由、DNS) |
无 | reused |
PublishPlane |
用一次存储,让新的路由表和目标成为新流使用的那一套 | 无 | 新的 Plane,epoch 加 1,负载均衡器探测 |
不报告 |
SwapHandler |
在入站已有的监听器上,用新 handler 服务它的新连接 | Listener::prepare_swap |
Listener::swap |
swapped |
Disrupt |
该变更会结束入站的存活连接 | 无:除非被允许,Actor::apply 中的中断检查会先拒绝它 |
没有它自己的工作 | restarted |
StopAccepting |
关闭入站在 bind 上的监听器:入站被移除或迁移 |
无 | Listener::stop 并移除它 |
不报告 |
CloseSessions |
关闭已不存在的入站 tag 的存活会话 | 无 | Sessions::close(Scope::Inbound) |
removed |
Drain |
让一个出站版本退出服务 | 无 | 按解析出的策略执行 drain |
drained |
Build(UserSet) 和 Reuse(UserSet) 本身不做任何工作。哪些用户表要变,由准备阶段逐个入站比较准入记录来决定(见准备阶段)。
Resource 与 OutboundId
Section titled “Resource 与 OutboundId”#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]pub struct OutboundId { pub tag: CompactString, pub version: u64,}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]pub enum Resource { Inbound(CompactString), Outbound(OutboundId), Balancer(CompactString), UserSet(CompactString), Route, // there is one Dns, // there is one}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]pub enum ResourceKind { Inbound, Outbound, Balancer, UserSet }| 值 | Display |
|---|---|
OutboundId { tag: "direct", version: 2 } |
direct@v2 |
Resource::Inbound("socks") |
inbound socks |
Resource::Outbound(direct@v2) |
outbound direct@v2 |
Resource::Balancer("lb") |
balancer lb |
Resource::UserSet("people") |
user set people |
Resource::Route |
route |
Resource::Dns |
dns |
ResourceKind |
inbound、outbound、balancer、user set |
OutboundId 取代了按指针地址识别身份的做法。出站的生命周期长于单次应用,一个 tag 可能有旧版本在新版本旁边排空,释放掉的内存也可能被它的后继复用,所以凡是必须区分两次构建的地方(UDP fan-out link 的子 link、流的 FlowEntry::outbound、排空步骤),都以 id 而不是 Arc 地址为键。负载均衡器的 Resource 只带 tag;它带版本的 id 用于它的 Target。
ApplyOptions
Section titled “ApplyOptions”#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct ApplyOptions { pub allow_disruptive: bool,}allow_disruptive 允许一次应用执行包含 Disrupt 步骤的计划。默认为 false。
ApplyReport
Section titled “ApplyReport”#[derive(Debug, Clone, Default, PartialEq, Eq)]pub struct ApplyReport { pub reused: Vec<Resource>, pub built: Vec<Resource>, pub swapped: Vec<CompactString>, pub drained: Vec<OutboundId>, pub rebound: Vec<CompactString>, pub restarted: Vec<CompactString>, pub removed: Vec<CompactString>,}| 字段 | 来源 | 含义 |
|---|---|---|
reused |
Step::Reuse |
原样保留:它们的 spec 没有变化 |
built |
Step::Build |
按 spec 构造:新增的,或有变化的;后一种情况下出站是其 tag 的一个新版本 |
swapped |
Step::SwapHandler |
保留监听器、用新 handler 服务新连接的入站 |
drained |
Step::Drain |
退出服务的出站版本;它们的存活流由排空策略处理 |
rebound |
tag 在旧 spec 中以另一个绑定存在的 Step::Bind |
监听器被关闭并重新绑定的入站 |
restarted |
Step::Disrupt |
变更结束了其存活连接的入站:这次应用被允许做出的中断性变更 |
removed |
Step::CloseSessions |
从 spec 中删去、会话已被关闭的入站 tag |
期望 spec 中的每个资源都按计划顺序出现在 reused 或 built 中。被移除的负载均衡器和被移除的用户集不出现在任何字段中;被移除的出站出现在 drained 中。因用户集变化而替换了用户表的入站记为 reused,该用户集记为 built。
ApplyError
Section titled “ApplyError”所有变体都在 supervisor/src/build/apply.rs 中用 thiserror 定义。规则类变体来自 validate,对应的规则见校验与 apply 错误。
| 变体 | 返回者 | Display |
|---|---|---|
DuplicateTag { kind, tag } |
plan(校验) |
duplicate {kind} tag {tag} |
UnknownReference { from, kind, tag } |
plan(校验) |
{from} references unknown {kind} {tag} |
EmptyBalancer { balancer } |
plan(校验) |
balancer {balancer} has no members |
UnprobeableMember { balancer, member } |
plan(校验) |
balancer {balancer}: outbound {member} has no upstream a TCP health probe can reach |
Invalid { resource, reason } |
plan(校验);set_users、upsert_user、remove_user |
{resource}: {reason} |
Disruptive { inbound, reason } |
Actor::apply 中的中断检查 |
inbound {inbound}: {reason}; this ends its live connections and needs allow_disruptive |
Build { resource, source } |
Actor::prepare、Actor::edit_users |
building {resource} failed: {source} |
Bind { inbound, bind, source } |
Actor::prepare |
inbound {inbound}: binding {bind} failed: {source} |
Stopped |
actor 退出之后,经由它的任何调用 | the supervisor has shut down |
Build 和 Bind 以底层的 io::Error 作为它们的 #[source]。{bind} 是 BindSpec 的 Display:TCP 为 127.0.0.1:1080,其余依次为 udp 127.0.0.1:8443、unix:<path>、tun <name>(或 tun auto)和 tun fd <n>。
下面是用锁定版本的 etemenanki-app 运行 --test 捕获的一条 Build 错误(--test 运行同样的准备过程,见 check),对应一个 cert_file 中没有 PEM 证书的 Hysteria 2 入站:
configuration invalid: building inbound hy2-in failed: hysteria2: the certificate file contains no certificatesconfiguration invalid: 是 app/src/main.rs → main 记录失败的 --test 时加的前缀;其余部分是 ApplyError 的 Display。
DrainPolicy
Section titled “DrainPolicy”#[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, // default Close pub drain: DrainPolicy, // default Keep}默认是 Keep。代码给出的理由是:TCP 流无法在中途迁移到后继版本,所以 supervisor 默认不关闭旧版本的流。UserRemovalPolicy(Keep、默认的 Close、CloseAfter(Duration))见用户与会话。两者在 Spec::policies 中各有一个 supervisor 级的默认值,出站(OutboundSpec::drain)或入站(InboundSpec::user_removal)可以覆盖它。
actor 在两次应用之间保存的状态
Section titled “actor 在两次应用之间保存的状态”调和的状态保存在私有的 Actor(supervisor/src/supervisor.rs)中,由 supervisor 唯一的那个任务持有:
| 字段 | 类型 | 保存的内容 |
|---|---|---|
state |
RunningState<U> |
最近一次提交的 spec,以及每个 tag 的最新版本 |
shared |
Shared |
每个 accept 循环共享的东西:plane 单元(PlaneCell)、Sessions 注册表、流 Tracker,以及持有每个 accept 循环和连接的 TaskTracker;负载均衡器探测和宽限期定时器也 spawn 在它上面 |
root |
CancellationToken |
每个监听器、探测和会话 token 的父 token |
dns、dns_target |
Option<Dns>、Option<Arc<Target>> |
正在运行的 DnsSpec 的解析器,以及存在 DNS 服务时它的内部目标 |
targets |
HashMap<CompactString, Arc<Target>> |
每个出站和负载均衡器 tag 的当前版本 |
probes |
HashMap<CompactString, CancellationToken> |
每个运行中负载均衡器的探测 token |
routes |
Option<Arc<CompiledRoutes>> |
正在运行的路由表 |
listeners |
Vec<Listener> |
每个已绑定的监听器,按 BindSpec 查找 |
admissions |
HashMap<CompactString, Admissions<U>> |
每个入站准入了谁,以入站 tag 为键:每个用户的 principal(身份主体)和凭据 |
keys |
UserKeys<U> |
每个被准入的用户 id 的 UserKey |
usage |
Arc<UsageBook<U>> |
用量账本,以及它的 key 所对应的 id;见按用户的用量计费 |
sampler |
Option<Sampler> |
统计采样器,直到 SupervisorBuilder::start 把它取走并 spawn |
background、background_stop |
Vec<JoinHandle<()>>、CancellationToken |
采样器和 usage sink 任务,以及在关闭结束时停止它们的 token |
epoch |
u64 |
已经发布过多少个 plane |
socket |
SocketOptions |
打开每个出站 socket 时遵循的 socket 策略;在 supervisor 的生命周期内固定不变 |
一次提交会替换 state、dns、dns_target、targets、probes、routes、listeners、admissions、keys 和 epoch,向 shared 的 plane 单元存入新值,并把新 key 的 id 发布到 usage。应用不会触碰 sampler、background、background_stop 和 socket;它们见 supervisor 概览。
准备阶段为一次应用构建、但还不对任何东西可见的内容:
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)>,}new_balancers 携带每个构建好的负载均衡器及其 probe_interval 和 probe_timeout。swaps 和 users 保存的是 Actor::listeners 中的位置。
flowchart TB
caller["Supervisor::apply_with"]
chan["经命令通道发送 Command::Apply"]
plan["plan:先 validate,再比较"]
gate{"有 Disrupt 步骤却没有 allow_disruptive?"}
prepare["Actor::prepare:构建、绑定、配对"]
commit["Actor::commit"]
ok["Ok:ApplyReport"]
err["Err:ApplyError,什么都没有改变"]
caller --> chan --> plan
plan -->|"违反某条规则"| err
plan --> gate
gate -->|"是"| err
gate -->|"否"| prepare
prepare -->|"构建或绑定失败"| err
prepare --> commit --> ok
所有会触发调和的入口最终都进入 Actor::apply:
impl<U: UserId> Supervisor<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>;}
impl<U: UserId> SupervisorBuilder<U> { pub fn options(self, options: ApplyOptions) -> Self; pub async fn start(self, spec: Spec<U>) -> Result<(Supervisor<U>, ApplyReport), ApplyError>;}旁边还有两个函数不经过 Actor::apply:epoch 读取 plane 单元,check 在它自己的 actor 上运行 plan 和 prepare。
impl<U: UserId> Supervisor<U> { pub fn epoch(&self) -> u64;}
pub async fn check<U: UserId>(spec: &Spec<U>) -> Result<(), ApplyError>;句柄通过一个私有的命令通道与 actor 通信。每条命令都带一个用于回复的 oneshot sender,Supervisor::ask 是唯一一个发送命令并等待回复的辅助函数:
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),}SupervisorBuilder::start先断言采样间隔不为零(the sample interval must not be zero)。然后它创建Actor,在 actor 任务存在之前,就在调用方的任务上用 builder 的options执行第一次应用。第一份 spec 被拒绝时返回它的ApplyError:actor 任务、采样器和 usage sink 都不会启动。只有第一次应用成功之后,start才 spawn 采样器,并在设置了 sink 时 spawn 用量推送任务,两者都在background_stop之下;随后它创建命令通道(mpsc::channel(16))并 spawnActor::run。从RunningState::empty()出发的计划从不包含Disrupt步骤,所以在此版本中options对第一次应用没有影响。apply(spec)等同于apply_with(spec, ApplyOptions::default())。两者都通过ask发送Command::Apply(spec, options, reply)并等待回复。Actor::apply依次运行plan(&self.state, &spec)、中断检查、prepare(&plan, &spec, true)和commit(plan, spec, prepared)。actor 一次只处理一条命令(Actor::run),所以两次应用永不交错,用户编辑、关闭和会话列举都要等整次应用(准备和提交)完成。这种串行来自唯一读取通道的那个任务,而不是来自锁。epoch()不经过 actor,直接从 plane 单元读取Plane::epoch。Plane::empty()的 epoch 为0,每次PublishPlane加一。句柄只在第一次应用成功后才存在,而第一次计划总会构建 DNS,因而总会发布,所以已启动的 supervisor 读到的值至少为1。
plan 按一个固定顺序产出步骤。提交阶段并不逐条重放这个列表:它先在一段固定的开头中发布 plane、撤销用户、启动监听器并切换 handler(见提交阶段),然后再遍历列表执行停止、关闭和排空,因此这三类步骤按图中的顺序执行。
flowchart LR dns["Dns"] --> out["每个出站"] --> bal["每个负载均衡器"] --> route["Route"] --> sets["每个用户集"] --> inb["每个入站:Bind、Build 或 Reuse"] inb --> pub["PublishPlane(plane 有变化时)"] --> swap["SwapHandler 与 Disrupt"] --> stop["StopAccepting"] --> close["CloseSessions"] --> drain["Drain"]
- 在同一类资源内部,Build 和 Reuse 步骤遵循期望 spec 中的顺序:出站按
Spec::outbounds的顺序,其余类推。 PublishPlane位于所有构建之后,所以发布的 plane 包含全部构建结果。Drain步骤排在最后,位于PublishPlane之后:只有取代旧版本的 plane 发布之后,旧版本才会退出服务。StopAccepting位于CloseSessions之前,所以被移除入站的监听器先被取消,然后才清扫它的会话。Sessions::open在注册表锁下检查监听器的停止 token,Sessions::close也在同一把锁下清扫,所以此刻正被 accept 的连接要么被清扫,要么根本不会打开会话(见用户与会话)。
plan 如何决定
Section titled “plan 如何决定”plan 先调用 validate(desired),出错就返回该错误。此后再没有任何东西会失败。它按上面的顺序遍历资源,一路维护 plane_changed、被重建的出站 tag 列表(rebuilt)、待定的 swaps 和待定的 drains。
没有运行中的 spec,或者 old.dns != desired.dns 时,dns_changed 为 true。它产出 Build(Dns) 或 Reuse(Dns);DNS spec 有变化时会设置 plane_changed,因为 plane 持有 DNS 服务的目标。
对每个期望的出站,prev 是运行中的 spec 里同 tag 的出站,running_version 是 running.versions[tag],只在 prev 存在时才查找。
| 情形 | 步骤 | 版本 |
|---|---|---|
prev == Some(outbound),并且不是(dns_changed 且该出站使用 DNS) |
Reuse(Outbound(tag@v)) |
不变 |
| spec 有变化,或 DNS 有变化且该出站不是 blackhole | Build(Outbound(tag@v+1)),并排入 Drain { tag@v, policy } |
next_version |
| 新 tag | Build(Outbound(tag@next)) |
next_version:从未出现过的 tag 为 1,否则比该 tag 上一个版本大一 |
每次构建都会设置 plane_changed,并把 tag 加入 rebuilt。除 blackhole 之外,每个出站都要解析域名,因此持有 DnsSpec 变化时会被替换的解析器;所以 DNS 变化会重建它们全部(uses_dns = !matches!(protocol, Blackhole))。
运行中的 spec 里有、期望的 spec 里没有的出站,得到 Drain { tag@v, policy } 并设置 plane_changed。它在 versions 中的条目保留。
drain_policy(tag) 在规划时一次性解析策略。它取带该 tag 的第一个出站,先在期望的 spec 中找,再在运行中的 spec 中找,使用它的 OutboundSpec::drain 覆盖值;该出站没有覆盖值时用 desired.policies.drain。因此,有变化的出站按新 spec 的覆盖值(或新的默认值)排空旧版本,被移除的出站则按它原来的覆盖值排空。
负载均衡器在以下条件全部满足时被复用:它的 BalancerSpec 未变、它有运行中的版本、它的成员都不在 rebuilt 中。负载均衡器持有其成员的目标,所以成员被重建时它也被重建,健康状态要重新探测。否则它得到 Build(Balancer(tag)),版本为 next_version(tag),并设置 plane_changed。运行中的 spec 里有、期望的 spec 里没有的负载均衡器会设置 plane_changed 并保留其版本。负载均衡器没有 Drain 步骤(见排空)。
出站和负载均衡器的 tag 共用一张版本表,正如它们在 spec 中共用一个命名空间。一个 tag 从出站变成负载均衡器(或反过来)时,版本号接着原来的计数。
没有运行中的 spec,或者 old.route != desired.route 时,route_changed 为 true。它产出 Build(Route) 或 Reuse(Route),有变化时设置 plane_changed。
运行中的 spec 持有同 tag 且相等的用户集时,期望的用户集为 Reuse(UserSet),否则为 Build(UserSet)。用户集从不设置 plane_changed,也从不影响入站的步骤。被移除的用户集不产生任何步骤。
入站按 BindSpec 而不是按 tag 匹配:绑定就是监听器的身份。
同一 BindSpec 上正在运行的入站 |
步骤 |
|---|---|
InboundSpec 相等 |
Reuse(Inbound(tag)) |
InboundSpec 不同 |
Build(Inbound(tag));排入:SwapHandler { tag },swap_disruption 给出原因时再排入 Disrupt |
| 无 | Bind { tag, bind }、Build(Inbound(tag));排入:变更替换了运行中 TUN 入站的设备时,排入一个 Disrupt(见中断性变更) |
InboundSpec 的相等比较覆盖每个字段(tag、bind、sniff、protocol、users、user_removal),所以在同一绑定上重命名入站或修改它的移除策略,都是一次 handler 切换。入站的变化从不设置 plane_changed。
接着,对运行中的入站:
- 没有任何期望入站使用的运行中
BindSpec,得到StopAccepting { tag, bind }; - 没有任何期望入站使用的运行中 tag,得到
CloseSessions { tag }。
因此,迁移的入站(tag 相同、绑定不同)是一个 Bind 加一个 StopAccepting,并保留它的会话,因为会话归属于 tag。在同一绑定上重命名的入站是一个 SwapHandler 加旧 tag 的 CloseSessions,并保留它的监听器。被移除的入站是一个 StopAccepting 加一个 CloseSessions。
最后,plan 在 plane_changed 时追加 PublishPlane,然后依次追加排队的切换和中断、停止、关闭和排空,并返回带有新 versions 的计划。
有两类变更天然会结束存活连接,plan 用 Step::Disrupt 标记它们:
| 变更 | 检测位置 | reason |
|---|---|---|
对保留设备(BindSpec)的 TUN 入站的任何修改,包括它的设置、sniff 或 tag |
swap_disruption:运行中的协议是 Tun |
the tun inbound's settings changed; its runtime restarts and its flows end |
同一绑定上 obfs 有变化的 Hysteria 2 入站 |
swap_disruption:两者都是 Hysteria2 且 old.obfs != new.obfs |
the hysteria2 obfuscation changed; connected clients cannot follow the new key |
| TUN 入站的设备有变化 | plan 中处理“没有运行中入站使用该绑定”的那个分支 |
the tun device changed; the old device and its flows end |
同一设备上 TUN 入站的每一次修改都是 Disrupt:只要运行中的协议是 Tun,swap_disruption 就返回一个原因。同一绑定上 Hysteria 2 入站的其他任何修改都是 SwapHandler,而不是 Disrupt。
Actor::apply 在准备任何东西之前检查 plan.disruptions().next()。没有 allow_disruptive 时,它针对按步骤顺序的第一个中断返回 ApplyError::Disruptive,例如:
inbound hy2-in: the hysteria2 obfuscation changed; connected clients cannot follow the new key; this ends its live connections and needs allow_disruptive各前端程序的选择:
- etemenanki-app 调用
Supervisor::apply,因此从不允许中断。app/src/instance.rs→Core::reload_with记录reload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply it,并保留正在运行的内容。所以 TUN 变更或 Hysteria 2 混淆变更需要重启才能生效。REST API 和 FFI 通过同一个Core重载(见etemenanki-app:运行、重载与关停)。 - katana 的
src/manager/node.rs→reconcile用ApplyOptions { allow_disruptive: true }应用每一份新的节点 spec,所以来自面板的 Hysteria 2 混淆变更会立即生效。见节点管理器。
每一行都是 supervisor/tests/unit/plan.rs 中的一个单元测试,从该行所说的 spec 出发:
| 变更 | 计划 |
|---|---|
第一次应用:一个准入用户集 people 的 SOCKS 入站,出站 direct 和 proxy,基于 proxy 的负载均衡器 lb |
Build(Dns)、Build(direct@v1)、Build(proxy@v1)、Build(Balancer lb)、Build(Route)、Build(UserSet people)、Bind socks 127.0.0.1:1080、Build(Inbound socks)、PublishPlane;版本 direct 1、proxy 1、lb 1 |
| 再次应用同一份 spec | 七个 Reuse 步骤,别无其他 |
direct 的地址族有变化 |
Build(direct@v2) … PublishPlane … Drain { direct@v1, Keep };路由和入站被复用 |
| 新增一条路由规则 | Build(Route)、PublishPlane;DNS、出站和入站被复用;没有排空 |
向 people 添加一个用户 |
Build(UserSet people)、Reuse(Inbound socks);没有绑定、停止、切换或发布 |
| 在同一绑定上把 SOCKS 入站改为 HTTP | Build(Inbound socks)、SwapHandler socks;没有绑定、停止、关闭或中断 |
| 把 SOCKS 入站迁移到端口 1081 | Bind socks 127.0.0.1:1081 在 Build(Inbound socks) 之前,StopAccepting socks 127.0.0.1:1080;没有切换,没有关闭 |
DNS spec 有变化,出站为 direct、hole(blackhole)和 proxy |
Build(Dns)、Build(direct@v2)、Reuse(hole@v1)、Build(proxy@v2)、PublishPlane,排空 direct@v1 和 proxy@v1 |
stateDiagram-v2 state "当前版本" as Current state "排空中" as Draining [*] --> Current: Build,随 plane 一起发布 Current --> Current: Reuse 保留同一个 Arc Current --> Draining: 重建或移除之后 Drain Draining --> [*]: Keep,drain 什么也不做 Draining --> [*]: Close,立即取消 closed token Draining --> [*]: CloseAfter,宽限期过后取消 closed token
- 一个 tag 的版本从 1 开始递增,永不复用。
next_version(tag)为running.versions[tag] + 1,从未应用过的 tag 为1。被移除的 tag 保留它的版本,所以再次加入时构建的是v+1。 - 被复用的出站保留它的版本和同一个
Arc<Target>,连同其 connector(连接器)持有的状态,例如 Hysteria 2 的 QUIC 连接或 WireGuard 隧道。构建出站只是构造它的 connector;在第一个流到来之前不会拨号。 - 在替换或移除旧版本的那次提交中,actor 不再持有旧版本的
Arc<Target>:self.targets在提交第 3 步被替换,旧表在提交返回时被丢弃。CloseAfter排空定时器会持有该Arc,直到它触发或被关闭流程结束。每个流保留它自己的 link 以及该版本的closedtoken,Close或CloseAfter会取消这个 token(见排空)。 - supervisor 自己创建两个目标,两者都使用空 tag,而 spec 中的 tag 不可能为空(
validate拒绝空 tag):一个是Plane::empty()的 blackhole@v0,它在第一次应用之前丢弃每个流;另一个是 DNS 服务的目标,其版本是构建时的epoch + 1。DNS 重建总会发布一个 plane,所以这个版本号正是首次携带它的那个 plane 的 epoch。
async fn prepare(&mut self, plan: &Plan, spec: &Spec<U>, bind: bool) -> Result<Prepared<U>, ApplyError>;Actor::prepare 完成所有可能失败的工作。它是 async 的,因为它要 await listener::bind,后者以异步方式打开 TUN 设备。它把结果构建在局部变量和一份暂存的用户 key 副本中,只在最后才组装 Prepared,所以在任何一步返回错误,都会丢弃到那时为止构建的一切,actor 保持不变。
| # | 步骤 | 做什么 | 错误 |
|---|---|---|---|
| 1 | DNS | 计划中有 Build(Dns) 或没有运行中的 DNS 时:build_dns(&spec.dns, plane, &socket),其中 plane 是指向 plane 单元的 Weak 引用(分流解析器中经代理解析的那一半通过它拨号);spec 中有 DNS 服务时,还创建一个新的内部目标 Target::outbound(internal_id(epoch + 1), Outbound::Dns(service))。否则克隆运行中的 Dns 和目标。 |
Build { resource: Dns } |
| 2 | 出站,按计划顺序 | Build(Outbound(id)):build_outbound(outbound, &dns, &socket),包装为 Target::outbound(id, built)。Reuse(Outbound(id)):克隆 self.targets[tag]。 |
Build { resource: Outbound(..) } |
| 3 | 负载均衡器,按计划顺序 | Build(Balancer(tag)):基于第 2 步得到的目标,为每个成员创建一个 Member::new(member, targets[member], probe_target(outbound)),然后 Balancer::new(members, strategy),包装为 Target::balancer(tag@plan.versions[tag], ..) 并排入 new_balancers。Reuse:克隆运行中的目标。 |
Build { resource: Balancer(tag) }:a balancer needs at least one outbound,校验已经排除了这种情况 |
| 4 | 路由 | 有 Build(Route) 或没有运行中的路由表时:compile_routes(&spec.route)(见 plane:逐流路由)。否则克隆运行中的 Arc。 |
Build { resource: Route } |
| 5 | 准入,逐个期望入站 | 在 self.keys 的副本上执行 admit(credential_kind(protocol), set, self.admissions.get(tag), &mut keys);被准入但没有可保留 principal 的用户会得到一个新的 principal:Principal::user(keys.key(id), id.name())。不再被准入的 principal 连同 removal_policy(spec, inbound) 一起排入 revoked。 |
无 |
| 6 | 入站,逐个期望入站 | 见下文。 | Build { resource: Inbound(tag) } 或 Bind |
第 3 步中的成员查找使用 expect("validated: members are outbounds") 和 expect("validated: members are probeable");第 2、3 步中的 spec 查找使用 expect("the plan builds the spec's outbounds") 和 expect("the plan builds the spec's balancers")。键缺失时会 panic 的普通下标访问也依赖同样的保证:复用出站时的 self.targets[&id.tag],构建负载均衡器时的 targets[member] 和 plan.versions[tag],复用负载均衡器时的 self.targets[tag],以及第 6 步中的 admissions[&inbound.tag]。校验和 plan 保证了所有这些。
第 6 步对每个期望入站执行如下操作,其中 running 是 bind 与该入站相同的那个监听器的位置:
- 计划中有
Build(Inbound(tag))。 用第 5 步得到的 principal 构建 handler:build_handler(inbound, admitted, circuits)。当该绑定上运行中的入站和期望的入站都是 Hysteria 2、且max_circuits相同(same_circuits)时,circuits是运行中 Hysteria 2 监听器的 circuit 配额(Listener::circuits()),这样存活的 circuit 会继续计入它;否则 handler 得到一个按新大小创建的新Semaphore。然后:- 有运行中的监听器时,
Listener::prepare_swap(handler)拒绝其他种类的 handler,检查 Hysteria 2 handler 的 Salamander 密钥(check_obfs,它把Salamander::new的错误转换为文本相同的InvalidInput错误),对 TUN 则复制监听器的设备(try_clone),并为切换时重启的运行时准备这份副本(TunInbound::prepare)。结果连同监听器的位置排入swaps; - 没有运行中的监听器且处于试运行(
bind == false)时,不再做别的; - 否则,
listener::bind(tag, bind).await打开 socket 或设备,Pending::pair(bound, handler)把它与 handler 配对,并构建服务所需的东西。这一对排入started。
- 有运行中的监听器时,
- 否则(入站被复用),当该绑定上有监听器在运行、且准入记录与运行中的不同(没有记录,或
same_admissions为 false)时,只重建用户表:user_table(inbound, admitted),然后Listener::prepare_users(table),后者拒绝其他种类的表。结果排入users。
listener::bind 和 Pending::pair 对各类绑定做什么(细节见监听器与服务循环):
| 绑定 | listener::bind |
Pending::pair |
|---|---|---|
| TCP | std::net::TcpListener::bind,设为非阻塞,包装给 tokio 使用 |
把监听器与 StreamInbound 配对 |
| Unix | 删除路径上已有的 socket 文件,视其为崩溃的上一次运行留下的残留;路径上是其他任何东西都拒绝,报 <path> exists and is not a socket(AlreadyExists);绑定后把文件的设备号和 inode 记录在 SocketFile 中 |
同 TCP |
| UDP(Hysteria 2) | std::net::UdpSocket::bind |
check_obfs,然后 Hy2Inbound::new 和 Hy2Inbound::open(socket),后者把 Salamander 密钥存入新的入站,并在该 socket 上打开 QUIC endpoint。endpoint 在运行之前不会从 socket 读取任何东西。 |
| TUN | 新建设备用 tun::open,使用外部提供的描述符用 tun::adopt,然后用 try_clone 得到一个 first 副本 |
TunInbound::prepare(first),它把副本注册到 tokio reactor(AsyncFd),并配置它的 IP 协议栈。在运行之前没有东西读取该设备。 |
绑定与 handler 的其他任何组合都会被拒绝,报 the handler is not of the kind its listener serves。
当两张表以相同顺序包含相同的用户,且每个用户的 principal 相同(Arc::ptr_eq)、凭据也相同时,same_admissions 成立,此时运行中的表已经表达了新表要表达的内容。只改了限速的用户,两者都会保留。
listener::bind 的错误变为 ApplyError::Bind { inbound, bind, source }。第 6 步中其他所有失败都是 ApplyError::Build { resource: Inbound(tag), source },例如:
- 无法解析的证书或私钥,例如
hysteria2: the certificate file contains no certificates(即上文的例子); - 来自
build_handler的inbound <tag>: protocol and bind do not match; - 来自
prepare_swap、prepare_users或Pending::pair的the handler is not of the kind its listener serves; - 被
check_obfs拒绝的 Salamander 密钥,或者Hy2Inbound::open、TunInbound::prepare的失败。
build_handler 还会为 Hysteria 2 的伪装(masquerade)调用 Masquerade::new,但校验会先调用它,所以被它拒绝的伪装永远到不了准备阶段。
准备阶段提前返回时会释放什么(check 丢弃它得到的 Prepared 时释放的也是这些):
- 第 6 步中绑定的 TCP 或 UDP socket 被关闭。Unix 监听器的 socket 文件也会被删除(由
SocketFile的Drop完成,且仅当该路径上仍是这个监听器创建的文件时;删除失败时以debug级别记录could not remove <path>: <e>)。新建的 TUN 设备的描述符被关闭; Pending::pair打开的 Hysteria 2 endpoint 随它的新Hy2Inbound一起被丢弃,这会结束 endpoint 的驱动任务。来自Pending::pair或prepare_swap的 TUNPreparedDevice被丢弃,它复制出的描述符被关闭;- 新的出站、负载均衡器、解析器和路由表被释放。它们都还没有启动任务或打开连接;
- 暂存的
UserKeys副本被丢弃,所以第 5 步分配的 key 都不会被发布。
中断检查在准备阶段之前运行,所以被它拒绝的 spec 什么都还没有构建。
check:试运行
Section titled “check:试运行”check(spec) 在一个新的 Actor(默认 socket 策略)中,以 RunningState::empty() 为起点运行 plan,然后运行 prepare(&plan, spec, false),并丢弃结果。它的文档注释说它“refuses what a start would, short of a port or device that cannot be bound”,即除了无法绑定的端口或设备之外,启动时会拒绝的它都会拒绝。它构建启动时会构建的一切:解析器、出站、负载均衡器、路由表、每个 handler 及其证书、每张用户表。它不绑定任何东西,也不配对,所以无法绑定的端口或设备,以及 Pending::pair 所做的检查(check_obfs、Hy2Inbound::open、TunInbound::prepare)都不会被触及。由于它从零开始规划,它永远不会遇到中断检查。etemenanki-app 的 --test(app/src/instance.rs → check)和 katana 的逐节点检查(src/runtime.rs → check_node)都调用它;见校验与 apply 错误。
Actor::commit(&mut self, plan: Plan, spec: Spec<U>, prepared: Prepared<U>) -> ApplyReport 不会失败。可能失败的事情都已在绑定、构建和配对时完成;剩下的 expect 和 unreachable! 调用只是陈述准备阶段已经保证的事(the key was checked by prepare_swap、prepare_swap pairs a swap with its listener's kind、prepare_users pairs a table with its listener's kind),第 4 步填充槽位时的普通下标 self.targets[tag] 也是如此。
| # | 步骤 | 为什么在这里 |
|---|---|---|
| 1 | commit_keys(keys):采用暂存的 key,并把新 key 的 id 发布到用量账本(UsageBook::names) |
会话绑定某个 key 之前,这个 key 必须已经发布:用量账本遇到它无法对应到 id 的 key 会 panic(… was bound before it was published) |
| 2 | publish_speed_limits(&spec):把每个用户的 speed_limit 交给 Tracker::set_speed_limits;属于多个用户集的用户取最小值;没有 UserKey 的用户(self.keys.get(id) 为 None)被跳过;set_speed_limits 设置它拥有的每个 pacer 的速率,不在表中的 key 设为不限速,并为每个有限速但还没有 pacer 的 key 插入一个 pacer |
只在提交时进行,所以被拒绝的应用永远不会改变正在生效的限速。存活的流在下一次读或写时感受到新的限速。 |
| 3 | 用准备好的表替换 self.targets,旧表保留为 old_targets |
第 11 步的排空需要旧版本 |
| 4 | 有 PublishPlane 时:epoch += 1;构建 Plane::new(routes, slots, dns_target, epoch),其中 slots 按 CompiledRoutes::targets() 的顺序为每个路由槽位保存一个目标(每个都是 self.targets[tag];Plane::new 用 debug_assert_eq! 检查两者长度一致);把它存入 plane 单元 |
一次存储同时换掉路由、出站和 DNS 服务,所以规则永远不会指向 plane 中没有的 tag |
| 5 | 有 PublishPlane 时:取消计划中不是 Reuse 的每个负载均衡器的探测 token;为每个新的负载均衡器创建根 token 的一个子 token,并在 task tracker 上用新的 dns.servers 解析器和 TcpDialer::new(socket) 执行 Balancer::spawn_probe |
不再位于 plane 中的负载均衡器停止探测;新的负载均衡器用其成员拨号时所用的解析器进行探测 |
| 6 | 保存 routes、dns 和 dns_target |
留给下一次应用的 Reuse |
| 7 | 对每一组被撤销的 principal 执行 Sessions::revoke(principals, policy) |
在存入任何新用户表之前(见撤销用户) |
| 8 | 对每个 started 对执行 Listener::start,各用根 token 的一个子 token,并追加到 listeners |
新监听器从第一个连接起就经由新 plane 路由 |
| 9 | 对每个 swaps 条目执行 Listener::swap |
准备阶段记录的位置仍然有效:第 8 步只做追加 |
| 10 | 对每个 users 条目执行 Listener::store_users |
在撤销之后 |
| 11 | 按顺序遍历计划中的步骤,填写报告;StopAccepting 停止并移除它的绑定上的监听器(swap_remove);CloseSessions 关闭该 tag 的会话;Drain 在 old_targets[tag] 恰好是那个版本时把旧目标交给 drain,并且无论如何都把该步骤的 id 推入 report.drained |
在所有切换和用户表存储之后,所以 swap_remove 不会移动仍要使用的位置;排空在发布之后 |
| 12 | self.admissions = admissions;self.state = RunningState { spec, versions: plan.versions } |
下一次规划从这里开始 |
| 13 | tracing::debug!("applied: {} built, {} reused, {} swapped, {} drained", …) |
日志 target 为 etemenanki_supervisor::supervisor |
没有 PublishPlane 时(只改了用户集或入站),plane、epoch 和探测都保持不动。Plane::empty() 是第一次提交之前的 plane,它通过基于空 tag 的 CompiledRoutes::single 把一切路由到它唯一的槽位,即 @v0 blackhole。
plane 只持有它的路由表指名的目标。只被负载均衡器指名的出站,要经由该负载均衡器的 Target 到达,后者持有准备阶段构建或复用的成员目标。
各监听器操作做什么
Section titled “各监听器操作做什么”这些操作位于 supervisor/src/system/listener.rs;它们驱动的 accept 循环见监听器与服务循环。
| 操作 | 流式监听器(TCP、Unix) | Hysteria 2 | TUN |
|---|---|---|---|
Listener::start |
在 task tracker 上 spawn run_stream_inbound,带一个装有 handler 的 watch 通道 |
spawn run_hysteria_inbound,带上 endpoint 和一个装有 tag 的 ArcSwap |
spawn run_tun_inbound,带一个无界的重启通道 |
Listener::swap |
把监听器的 tag 设为新入站的 tag,然后把新的 Arc<StreamInbound> 用 send_replace 送入 watch 通道;切换之后 accept 的连接由它服务 |
密钥有变化时调用 set_obfs;设置监听器的 tag,并把新 tag 存入 run_hysteria_inbound 每个连接读取一次的 ArcSwap,因此重命名入站的新连接会记在新 tag 之下;用 handler 的 circuit 配额替换保留的配额;然后在运行中的 Hy2Inbound 上调用 set_connection_config、set_quic 和 set_max_connections |
设置监听器的 tag,然后把新 handler 和准备好的副本送上重启通道;发送失败时,spawn_tun 用它们启动一个运行时任务,并用它的 sender 取代旧的 |
Listener::store_users |
存入新的用户表(UserTable::store_into);见用户与会话 |
set_authenticator |
无:TUN 设备不准入任何用户 |
Listener::stop |
取消监听器的 token,accept 循环随之停止。会话运行在根 token 的子 token 之下,而不是这个 token 之下。 | 取消监听器的 token;endpoint 如何收尾见监听器与服务循环 | 取消监听器的 token;运行时任务的每一轮运行都持有它的一个子 token |
每个启动的监听器都以 info 级别记录 inbound <tag> listening on <bind>(target 为 etemenanki_supervisor::system::listener)。在准备阶段绑定时,新建的 TUN 设备还会记录 inbound <tag> owns tun device <name>,外部提供的设备则记录 inbound <tag> serves a supplied tun device。
一个 Drain 步骤最终落到一次调用:
fn drain(target: Arc<Target>, policy: DrainPolicy, root: &CancellationToken, tracker: &TaskTracker);| 策略 | drain 对旧版本的流做什么 |
|---|---|
Keep(默认) |
什么也不做:supervisor 不关闭它们 |
Close |
Target::close_flows():立即关闭它们 |
CloseAfter(grace) |
在 tracker 上 spawn 一个任务,在 root.cancelled() 和 sleep(grace) 之间 select;sleep 结束后调用 close_flows(),因此在 grace 之后关闭它们 |
无论哪种策略,新流都遵循已发布的 plane:重建之后走该 tag 的新版本,移除之后则去往新路由表指定的地方。
关闭的工作方式(supervisor/src/topology/plane.rs):
- 每个
Target拥有一个closed: CancellationToken。close_flows()取消它。 Target::connect_stream和Target::connect_datagram先解析目标(负载均衡器会选出一个成员),再把打开的 link 连同解析出的目标的closedtoken 一起包装进Guarded<T>。如果解析出的目标本身又是负载均衡器,调用会以balancer member <id> is itself a balancer失败,而校验已经排除了这种情况。Guarded在每次poll_read、poll_write、poll_flush、poll_send_to和poll_recv_from之前检查该 token,包括已经在等待的接收,并返回文本为the outbound this flow was opened on was drained的io::ErrorKind::ConnectionAborted。poll_shutdown不经检查直接透传。中继该流的运行时看到一个出站错误,并结束这个流。- 经过负载均衡器的流是在为它选出的成员上打开的,所以它携带的是该成员的 token,并由该成员的排空关闭。因此负载均衡器没有
Drain步骤:成员都被复用的负载均衡器即使被重建,也不会关闭这些流;被重建的成员则排空它自己的流。 - DNS 服务的内部目标同样没有
Drain步骤,所以 supervisor 不会关闭已经交给旧服务的流。 - 尚未触发的
CloseAfter排空定时器不会让shutdown等待超过shutdown自己的grace:定时器运行在 task tracker 上,而shutdown对 task tracker 最多等待它的grace;定时器还在根 token 上 select,shutdown随后会取消根 token(shutdown_ends_pending_grace_timers)。
排空策略在规划时解析(见出站)并固定在步骤中,所以之后的应用无法改变已排空版本的策略。
用户的变化要么来自一次应用(UserSet 有变化,或者入站改为准入另一个用户集、读取另一种凭据类型),要么来自下文的用户编辑路径。无论哪种,build/users.rs → admit 都会把入站运行中的准入记录与新的进行比较:
- 仍被准入、且仍持有当初据以被准入的那个凭据的用户,保留其
Arc<Principal>,即使入站现在读取另一种凭据类型也是如此,因此他们的存活会话不受影响; - 不再被准入、或之前的凭据已改变或消失的用户,失去其旧 principal,旧 principal 在
revoked中返回。如果他们仍被准入,就得到一个新的 principal:Principal::user(keys.key(id), id.name())。
UserKeys::key 返回用户 id 已有的 key,或者分配下一个。key 从 1 开始递增,所以 key n 对应 ids[n - 1](build/users.rs → id_of)。
每个入站被撤销的 principal 连同该入站的策略一起排队:inbound.user_removal,未设置时为 spec.policies.user_removal。准入记录以入站 tag 为键,所以重命名的入站从全新的 principal 开始;它旧 tag 的会话由 CloseSessions 关闭,而不是被撤销。
在提交阶段,撤销在存入任何新用户表之前运行,也在新监听器启动、handler 切换之前运行:
sequenceDiagram participant C as Actor 的 commit participant P as 旧 Principal participant H as 基于旧用户表的握手 participant S as 它的 Session C->>P: Sessions::revoke 将其标记为已撤销 C->>S: 绑定到它的存活会话按策略处理 H->>P: 以旧 principal 的身份认证 H->>S: 第一个流调用 Session::admit S-->>H: 拒绝,principal 已被撤销 C->>C: store_users 装入新表
Sessions::revoke(principals, policy) 获取注册表锁,然后对每个 principal:
- 按策略把它标记为已撤销(
Principal::revoke;只有第一次撤销有效); - 在注册表的
by_user索引中查找其UserKey的会话,只保留所绑定的 principal 正是这同一个Arc(Arc::ptr_eq)的会话。同一用户在另一个入站上的会话持有另一个 principal,不受影响; - 对其中每个会话执行策略(
enforce):Keep什么也不做,Close取消会话的 token,CloseAfter(grace)在 task tracker 上 spawnclose_after,它在会话的 token 和sleep(grace)之间 select,sleep 结束后取消该 token。
第一个流出示的 principal 已被撤销的会话会被拒绝:第一个流绑定会话时,Session::admit 在同一把注册表锁下检查 principal,并且在任何策略下都以 the user was removed before the session opened a flow 拒绝它,因为移除策略只宽待移除时已经绑定到该 principal 的会话。在 store_users 之前撤销,使这条规则也覆盖那些针对即将被替换的旧表完成认证的握手。Principal、Session 和各种策略的机制见用户与会话。
用户编辑路径
Section titled “用户编辑路径”set_users、upsert_user 和 remove_user 不经过规划。它们发送 Command::Users,actor 运行 edit_users,即针对单个用户集的一套更小的准备和提交:
- 克隆运行中的 spec;没有时回复
ApplyError::Stopped。编辑指定的用户集:Set替换它的全部用户,Upsert插入或替换一个用户,Remove移除一个用户。用户集不存在时为ApplyError::Invalid { resource: UserSet(tag), reason: "no such user set" }(user set <tag>: no such user set)。Set和Upsert总被视为一次变更,即使结果完全相同。Remove一个不在集合中的用户会返回Ok(false),不做任何改变。 - 准备:对
users指向该用户集的每个入站,依次执行validate_admission(inbound, set)(对 TUN 入站立即返回Ok,因为它没有凭据类型);在 key 的暂存副本上执行admit;取入站绑定上的监听器(expect("every running inbound has its listener"));user_table;Listener::prepare_users。每个这样的入站的用户表都会重建,不做same_admissions检查。任何失败都在存储任何东西之前返回:validate_admission的错误原样返回,user_table或prepare_users的错误作为ApplyError::Build { resource: Inbound(tag) }返回。这一步会用expect("found above")再次查找被编辑的用户集。 - 提交:
commit_keys、publish_speed_limits,按各入站的策略撤销其被移除的 principal,然后存入每张用户表及其准入记录,最后把编辑后的 spec 记为运行中的 spec。
plane、监听器、handler 和版本都不受影响,也不生成 ApplyReport。remove_user 返回该用户是否在集合中;set_users 和 upsert_user 返回 ()。
update 与 update_with
Section titled “update 与 update_with”update(edit) 等同于 update_with(edit, ApplyOptions::default())。它发送 Command::Update(Box::new(edit), options, reply),其中装箱的闭包是一个 Edit<U>(见一次应用)。actor 克隆它运行中的 spec,在 actor 任务内部对副本调用 edit,然后完全按 apply_with 的方式应用结果。由于编辑在 actor 上运行,它看到的总是上一条命令留下的 spec,两个调用方之间不存在读-改-写竞争。没有运行中的 spec 时它回复 ApplyError::Stopped,而已启动的 supervisor 总有运行中的 spec。在此版本中,仓库内的前端程序都没有调用它。
一次应用保留什么
Section titled “一次应用保留什么”| 资源 | 保留条件 | 保留下来的内容 |
|---|---|---|
| 监听器 | 它的 BindSpec 不变(Reuse 或 SwapHandler) |
socket 或设备及其 accept 循环;对流式监听器,还有 run_stream_inbound 只创建一次的 MAX_HANDSHAKES_PER_INBOUND(204,800)和 MAX_LIVE_CONNECTIONS_PER_INBOUND(4,194,304)信号量;对 Hysteria 2,还有 QUIC endpoint,以及 max_circuits 相同时的 circuit 配额 |
| 出站 | 它的 spec 不变;除非它是 blackhole,DNS spec 也不变 | 同一个 Arc,以及其 connector 持有的状态,例如 QUIC 连接或隧道 |
| DNS 解析器 | DnsSpec 不变 |
解析器及其共享的应答缓存,以及 DNS 服务的目标 |
| 负载均衡器 | 它的 spec 和每个成员都不变 | 成员的健康状态和运行中的探测 |
| 路由表 | RouteSpec 不变 |
编译好的路由表 |
| 用户的 principal | 用户仍以同一个凭据被准入 | 他们的存活会话;限速变化也会保留这些会话 |
| plane | DNS、出站、负载均衡器和路由都没有变化 | 已发布的 Plane 和 epoch |
会话运行在根 token 的子 token 之下,而不是其监听器 token 的子 token 之下,所以停止监听器或切换它的 handler 都不会取消会话。调和通过两种途径关闭会话:对离开 spec 的入站 tag 执行 CloseSessions;通过 Sessions::revoke,对绑定到已撤销 principal 的会话执行移除策略。Disrupt 步骤标记的变更,其原因说明它会结束入站的存活连接(见中断性变更)。
策略(Spec::policies)没有对应的步骤。新的默认值从带来它们的那次应用的撤销和排空开始生效,此后一直有效。
路由变化不会迁移已建立的流:connector 每次 connect 读取一次 plane(supervisor/src/connector.rs → AppConnector::connect),UDP fan-out link 每发送一个数据包读取一次(FanOutLink::poll_send_to),所以新连接、每个新的 mux 子流,以及 fan-out link 的每个 UDP 数据包,都遵循它们被路由时的当前 plane(见 plane:逐流路由)。
| 不变量 | 由谁保证 | 由谁固定 |
|---|---|---|
| 规划没有副作用 | plan 只接受引用,不做 I/O |
supervisor/tests/unit/plan.rs 中的每个测试都手工构造 spec,不绑定任何东西 |
| 在提交之前任何一处被拒绝的 spec 都不改变任何东西 | plan 中的校验;准备之前的中断检查;准备阶段构建到局部变量和暂存的 key 中;限速和 key 只在提交时发布 |
a_refused_spec_changes_nothing、an_invalid_spec_is_refused_without_a_plan |
| 为被拒绝的 spec 绑定的监听器会被释放 | 准备阶段的局部变量拥有绑定得到的句柄;提前返回时丢弃它 | a_refused_spec_changes_nothing |
| tag 的版本永不复用 | versions 保留已移除的 tag;next_version 递增 |
each_rebuild_of_a_tag_takes_the_next_version、a_tag_added_back_takes_a_version_it_never_had、a_removed_outbound_drains_under_its_own_policy |
| 只有取代旧版本的 plane 发布之后,旧版本才会被排空 | Drain 步骤位于 PublishPlane 之后 |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
| 只排空步骤所指的那个版本 | 提交中的 target.id() == outbound 检查;在此版本中 old_targets[tag] 总是步骤所指的版本,所以这个检查是防御性的 |
无 |
| 未变的出站保留其状态 | Reuse 克隆运行中的 Arc |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
| DNS 变化会传达到每个持有解析器的对象 | 按 uses_dns 重建每个非 blackhole 出站 |
a_dns_change_rebuilds_every_outbound_but_a_blackhole |
| 负载均衡器随其成员一起重建,否则保留 | members_kept |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is |
| 移除负载均衡器会停止指向它的路由 | 移除时设置 plane_changed |
removing_a_balancer_alone_republishes_the_plane |
| 用户变化从不触及 plane 或监听器 | 用户集从不设置 plane_changed;用户表原地换入 |
a_user_set_change_alone_touches_neither_the_plane_nor_the_inbounds、a_speed_limit_change_keeps_the_users_sessions |
| 未变的绑定保留其 socket | 入站按 BindSpec 匹配;SwapHandler |
a_protocol_change_on_the_same_bind_swaps_the_handler_and_keeps_the_listener、a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
| 迁移的入站保留其会话 | 会话归属于 tag | a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions |
| 被移除或被重命名的入站的会话会被关闭 | 对每个消失的 tag 执行 CloseSessions,位于 StopAccepting 之后 |
a_removed_inbound_stops_accepting_and_closes_its_sessions、an_inbound_renamed_on_the_same_bind_swaps_and_closes_the_old_tags_sessions、removing_an_inbound_closes_its_sessions、a_stopped_listener_opens_no_session |
| 中断性变更需要明确同意 | Actor::apply 中的中断检查 |
a_hysteria2_obfuscation_change_needs_allow_disruptive、a_tun_device_change_needs_allow_disruptive |
| 第一个流出示已撤销 principal 的会话会被拒绝 | Session::admit 在注册表锁下检查 principal;撤销在 store_users 之前运行 |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy、removing_a_user_closes_only_their_sessions_and_refuses_them_after |
| key 在会话能绑定它之前就已发布 | 在提交和 edit_users 中都最先执行 commit_keys |
无;违反时用量账本会 panic |
| epoch 只在发布 plane 时变化 | 只在 PublishPlane 下执行 epoch += 1 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow(epoch + 1)、a_refused_spec_changes_nothing(不变) |
| 宽限期定时器随 supervisor 一起结束 | 排空定时器在根 token 上 select;移除定时器在会话的 token(根 token 的子 token)上 select | shutdown_ends_pending_grace_timers、cancelling_the_root_cancels_every_session |
| 准备阶段记录的位置在提交时仍然有效 | 新监听器只做追加;swap_remove 只在步骤遍历中运行,位于所有切换和存储之后 |
无 |
失败路径与取消
Section titled “失败路径与取消”| 位置 | 结果 | 之后的状态 |
|---|---|---|
plan |
一个校验类 ApplyError |
不变 |
| 中断检查 | 指出第一个中断的 ApplyError::Disruptive |
不变;什么都没有构建 |
| 准备 | 第一个失败对应的 ApplyError::Build 或 ApplyError::Bind,按准备阶段表格中的顺序 |
不变;已构建的东西被丢弃并释放 |
| 提交 | 不会失败 | 新的 spec 开始运行 |
| actor 已退出 | apply、update 和用户方法返回 ApplyError::Stopped;close 返回 0;sessions 返回空列表;shutdown 直接返回 |
supervisor 已关闭 |
- Stopped。
Supervisor::ask把已关闭的命令通道和被丢弃的回复都映射为ApplyError::Stopped。actor 在处理Command::Shutdown之后返回;最后一个Supervisor句柄被丢弃时,它以零宽限期自行关闭。 - 取消调用。 在
applyfuture 等待通道容量时丢弃它,什么也不会发送。命令一旦入队,无论调用方怎样,actor 都会把这次应用执行到底;调用方已经放弃时,回复随之被丢弃。 - 启动。 第一份 spec 被拒绝时,
SupervisorBuilder::start以该ApplyError失败;actor 任务、采样器和 usage sink 都不会启动,准备阶段的局部变量按上文所述被丢弃。采样间隔为零时,start在第一次应用之前 panic。 - 没有运行中的 spec。 actor 没有运行中的 spec 时,
Command::Update和edit_users回复ApplyError::Stopped。已启动的 supervisor 总有一份,所以这只是一道防护,不是调用方能走到的路径。
前端程序如何包装结果:
- etemenanki-app 的
Core::reload_with对大多数拒绝记录reload: <error>; keeping the running config,对Disruptive记录上文那条专门的日志,成功时记录config reloaded: <summary>。summary 按 built、swapped、drained、rebound、restarted、removed、reused 的顺序列出报告中非空的分组,每组写成label [item, item],各组以;连接;所有分组都为空时为nothing to run。见etemenanki-app:运行、重载与关停。 - katana 为对应的节点记录应用的错误;见节点管理器。
| 项目 | 值 | 定义位置 |
|---|---|---|
| 命令通道 | mpsc::channel(16) |
SupervisorBuilder::start |
| 同时进行的应用 | 一个:actor 串行处理每条命令 | Actor::run |
| plane epoch | u64,Plane::empty() 为 0,每次发布 +1 |
Actor::commit |
| 出站和负载均衡器版本 | 每个 tag 一个 u64,从 1 开始 |
plan → next_version |
versions 表 |
每个曾经应用过的出站或负载均衡器 tag 一个条目;从不裁剪 | RunningState、Plan |
| 内部目标 | 空 tag;空 plane 为 @v0,DNS 服务为 @v<epoch + 1> |
Plane::empty、internal_id |
| TUN 重启通道 | mpsc::unbounded_channel,只由 actor 写入,每次切换写入一次 |
system/listener.rs → spawn_tun |
| 流式 handler 通道 | watch::channel,一个值 |
Listener::start |
| 资源查找 | 按 tag 或按绑定对 spec 中的各个 vector 线性扫描 | plan、Actor::prepare、Actor::commit |
supervisor/tests/unit/plan.rs 单独驱动 plan。它的辅助函数构建一个 127.0.0.1:1080 上的 SOCKS 入站和一个 freedom 出站 direct(spec()),以及一个 full_spec():在此基础上增加一个指向 192.0.2.1 的 SOCKS 客户端出站 proxy、一个基于它的 failover 负载均衡器 lb(每 30 秒探测一次,超时 5 秒),以及一个由该入站准入的用户集 people。TUN 入站只是一个 DeviceSpec 值,所以不会打开任何东西。
| 测试 | 固定的行为 |
|---|---|
the_first_plan_binds_and_builds_everything_then_publishes |
第一次应用的确切步骤列表和版本 |
an_unchanged_spec_reuses_everything_and_publishes_nothing |
七个 Reuse 步骤,每个资源一个;版本不变 |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
构建 v2,在 PublishPlane 之后以 Keep 排空 v1;路由和入站被复用 |
a_drain_takes_the_supervisor_default_unless_the_outbound_overrides_it |
policies.drain = Close 生效;出站的 CloseAfter(5 s) 覆盖值优先 |
each_rebuild_of_a_tag_takes_the_next_version |
两次重建得到 v3 并排空 v2 |
a_removed_outbound_drains_under_its_own_policy |
按它原来的 Close 覆盖值排空;plane 重新发布;版本保留 |
a_tag_added_back_takes_a_version_it_never_had |
移除之后再加回,得到 v2 |
a_dns_change_rebuilds_every_outbound_but_a_blackhole |
freedom 和 proxy 被重建并排空;blackhole 被复用 |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is |
非成员重建时保留 lb@v1;成员重建时构建 lb@v2 并发布 |
removing_a_balancer_alone_republishes_the_plane |
发布,不排空,版本保留 |
a_route_change_alone_rebuilds_only_the_route_table |
只有 Build(Route) 和一次发布 |
a_user_set_change_alone_touches_neither_the_plane_nor_the_inbounds |
Build(UserSet)、Reuse(Inbound),没有发布、绑定、停止或切换 |
a_protocol_change_on_the_same_bind_swaps_the_handler_and_keeps_the_listener |
SwapHandler,没有绑定、停止、关闭或中断 |
a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions |
Bind 在 Build 之前,旧绑定的 StopAccepting,没有关闭 |
a_removed_inbound_stops_accepting_and_closes_its_sessions |
计划以 StopAccepting 加 CloseSessions 结尾 |
an_inbound_renamed_on_the_same_bind_swaps_and_closes_the_old_tags_sessions |
新 tag 的 SwapHandler,旧 tag 的 CloseSessions,没有绑定或停止 |
a_tun_device_change_opens_a_new_device_and_disrupts |
变化的 DeviceSpec 是一个 Bind、一个 StopAccepting 和一个 Disrupt,没有切换 |
a_tun_settings_change_on_the_same_device_swaps_and_disrupts |
变化的 TunSpec 是一个 SwapHandler 和一个 Disrupt |
a_hysteria2_obfs_change_swaps_and_disrupts |
新的 Salamander 密钥是一个 SwapHandler 和一个 Disrupt |
a_hysteria2_masquerade_change_swaps_without_disrupting |
伪装变化是一个 SwapHandler,没有 Disrupt |
an_invalid_spec_is_refused_without_a_plan |
未知的默认出站返回 UnknownReference { from: Route } |
supervisor/tests/hot_swap.rs 借助 supervisor/tests/support/mod.rs 中的辅助函数,在回环 socket 上运行真实的 supervisor(base_spec():出站 direct 和 block,默认路由 direct,系统 DNS;SOON = 5 s,QUIET = 400 ms):
| 测试 | 固定的行为 |
|---|---|
an_established_connection_survives_an_apply_that_keeps_its_inbound |
增加一个出站和一条规则的应用复用了入站;已打开的 SOCKS 连接继续回显 |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
epoch 增加一;存活 SOCKS UDP 关联的下一个数据包走新路由;已打开的 TCP 流不变;新的 TCP 流走新路由 |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
在同一个存活的 VLESS mux 会话上,变更之后打开的子流走新路由,之前的子流保持它自己的路由 |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close |
以 freedom 出站 direct 为例:在 Keep 下,构建 v2 的那次应用之后,在 direct@v1 上打开的流仍在回显;在 Close 下,v2 上的流被关闭,v1 上的流仍在回显(v1 was not drained again),新流使用 v3 |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
被复用的 Hysteria 2 出站不会打开第二个 QUIC 连接;被重建的出站重新拨号,而旧连接上的流继续进行 |
a_hysteria2_obfuscation_change_needs_allow_disruptive |
以 Disruptive 被拒绝;允许之后,报告 restarted 且没有 rebound |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
同一端口上 SOCKS 改为 HTTP:swapped,没有 rebound;已打开的 SOCKS 连接继续中继;新连接说 HTTP |
a_refused_spec_changes_nothing |
在一次 UnknownReference 拒绝和一次 Bind 拒绝之后:epoch 相同,路由和用户仍是旧的,准备期间绑定的监听器已被释放 |
removing_a_user_closes_only_their_sessions_and_refuses_them_after |
在 Close 和 Keep 下:只有被移除用户的会话被关闭(或者一个都不关闭),该用户的下一次 SOCKS 握手失败,close(Selector::User) 关闭另一个用户的会话 |
a_speed_limit_change_keeps_the_users_sessions |
以新限速调用 set_users 会保留会话及其 id |
removing_an_inbound_closes_its_sessions |
removed 列出该 tag;它的连接被关闭,另一个入站的连接不受影响;它的端口拒绝连接 |
a_tun_device_change_needs_allow_disruptive |
没有该选项时,设备变化和设置变化都被拒绝。需要 CAP_NET_ADMIN;无法创建设备时跳过 |
shutdown_ends_pending_grace_timers |
一小时的移除宽限期和一小时的排空宽限期都不会拖住 shutdown(Duration::ZERO) |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
客户端卡在 QUIC 握手中时,shutdown(Duration::ZERO) 在 SOON(5 s)内返回;它不会等到握手的空闲超时(见 supervisor 概览) |
这条路径上的其他测试:
| 文件 | 测试 | 固定的行为 |
|---|---|---|
supervisor/tests/unit/plane.rs |
closing_a_targets_flows_aborts_its_open_streams |
close_flows 之后,受保护 stream 的读、写和 flush 都会失败 |
supervisor/tests/unit/plane.rs |
closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive |
已在等待的接收失败,下一次发送也失败 |
supervisor/tests/unit/plane.rs |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers |
关闭负载均衡器的目标不影响 stream;关闭成员的目标则结束它 |
supervisor/tests/unit/session.rs |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy |
撤销必须先于用户表存储的原因 |
supervisor/tests/unit/session.rs |
revoking_under_close_ends_only_the_sessions_of_that_principal |
只有绑定到被撤销 principal 的会话被关闭;同一用户在另一个入站上的会话(持有另一个 principal)以及另一个用户的会话都保持打开 |
supervisor/tests/unit/session.rs |
revoking_under_keep_ends_no_session |
已绑定的会话保持打开,仍可准入流;principal 在 Keep 下被标记为已撤销 |
supervisor/tests/unit/session.rs |
revoking_under_close_after_ends_the_session_once_the_grace_passes |
宽限期为 10 s 时,会话在到期前 1 s 仍然打开,过期后被关闭 |
supervisor/tests/unit/session.rs |
cancelling_the_root_cancels_every_session |
无论是否已绑定,每个会话都是根 token 的子 token |
supervisor/tests/unit/session.rs |
a_stopped_listener_opens_no_session |
StopAccepting 必须先于 CloseSessions 的原因 |
app/tests/integration/e2e_reload.rs |
an_open_connection_keeps_transferring_across_a_reload |
经由 app:一次把新连接送往 blackhole 的重载之后,已打开的连接保持原来的路由 |
app/tests/integration/e2e_reload.rs |
a_reload_that_removes_a_user_closes_only_their_connections |
经由 app:从配置中移除一个用户,只关闭该用户的连接 |
app/tests/integration/e2e_hysteria_inbound.rs |
a_reload_keeps_the_inbound_serving_on_its_udp_port |
经由 app:绑定不变而设置有变化的 Hysteria 2 入站继续在它的 UDP 端口上服务 |
update 和 update_with、DrainPolicy::CloseAfter 定时器的触发、edit_users 的错误(no such user set),以及 Listener::swap 中 spawn_tun 的回退路径,都没有专门的测试。修改这些路径时请补上。