跳转到内容

规划并应用变更

源码文件:41 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
  • Etemenanki/supervisor/src/topology/spec_plan/plan.rs
  • Etemenanki/supervisor/src/topology/spec_plan/mod.rs
  • Etemenanki/supervisor/src/topology/spec_plan/inbound.rs
  • Etemenanki/supervisor/src/topology/spec_plan/outbound.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/topology/outbound/udp_fanout.rs
  • Etemenanki/supervisor/src/build/apply.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/policy.rs
  • Etemenanki/supervisor/src/topology/plane.rs
  • Etemenanki/supervisor/src/topology/router.rs
  • Etemenanki/supervisor/src/topology/inbound/mod.rs
  • Etemenanki/supervisor/src/topology/balancer.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/system/listener.rs
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/build/users.rs
  • Etemenanki/supervisor/src/build/inbound.rs
  • Etemenanki/supervisor/src/build/outbound.rs
  • Etemenanki/supervisor/src/build/dns.rs
  • Etemenanki/supervisor/src/build/route.rs
  • Etemenanki/supervisor/src/build/validate.rs
  • Etemenanki/supervisor/src/entity/session.rs
  • Etemenanki/supervisor/src/track/mod.rs
  • Etemenanki/supervisor/src/lib.rs
  • Etemenanki/protocols/src/dns/mod.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/supervisor/tests/unit/plan.rs
  • Etemenanki/supervisor/tests/hot_swap.rs
  • Etemenanki/supervisor/tests/support/mod.rs
  • Etemenanki/supervisor/tests/unit/plane.rs
  • Etemenanki/supervisor/tests/unit/session.rs
  • Etemenanki/app/src/main.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/api.rs
  • Etemenanki/ffi/src/proxy.rs
  • Etemenanki/app/tests/integration/e2e_reload.rs
  • Etemenanki/app/tests/integration/e2e_hysteria_inbound.rs
  • katana/src/manager/node.rs
  • katana/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:逐流路由。
supervisor/src/topology/spec_plan/plan.rs
#[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。
supervisor/src/topology/spec_plan/plan.rs
#[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) 本身不做任何工作。哪些用户表要变,由准备阶段逐个入站比较准入记录来决定(见准备阶段)。

supervisor/src/entity/id.rs
#[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。

supervisor/src/supervisor.rs
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ApplyOptions {
pub allow_disruptive: bool,
}

allow_disruptive 允许一次应用执行包含 Disrupt 步骤的计划。默认为 false。

supervisor/src/build/apply.rs
#[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。

所有变体都在 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 certificates

configuration invalid: 是 app/src/main.rs → main 记录失败的 --test 时加的前缀;其余部分是 ApplyError 的 Display。

supervisor/src/policy.rs
#[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 概览。

准备阶段为一次应用构建、但还不对任何东西可见的内容:

supervisor/src/supervisor.rs
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:

supervisor/src/supervisor.rs
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。

supervisor/src/supervisor.rs
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 是唯一一个发送命令并等待回复的辅助函数:

supervisor/src/supervisor.rs
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))并 spawn Actor::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。

update 和 check 分别见下文和准备阶段一节。

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 先调用 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 以及该版本的 closed token,Close 或 CloseAfter 会取消这个 token(见排空)。
  • supervisor 自己创建两个目标,两者都使用空 tag,而 spec 中的 tag 不可能为空(validate 拒绝空 tag):一个是 Plane::empty() 的 blackhole @v0,它在第一次应用之前丢弃每个流;另一个是 DNS 服务的目标,其版本是构建时的 epoch + 1。DNS 重建总会发布一个 plane,所以这个版本号正是首次携带它的那个 plane 的 epoch。
supervisor/src/supervisor.rs
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 的 TUN PreparedDevice 被丢弃,它复制出的描述符被关闭;
  • 新的出站、负载均衡器、解析器和路由表被释放。它们都还没有启动任务或打开连接;
  • 暂存的 UserKeys 副本被丢弃,所以第 5 步分配的 key 都不会被发布。

中断检查在准备阶段之前运行,所以被它拒绝的 spec 什么都还没有构建。

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 到达,后者持有准备阶段构建或复用的成员目标。

这些操作位于 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 步骤最终落到一次调用:

supervisor/src/supervisor.rs
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 连同解析出的目标的 closed token 一起包装进 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:

  1. 按策略把它标记为已撤销(Principal::revoke;只有第一次撤销有效);
  2. 在注册表的 by_user 索引中查找其 UserKey 的会话,只保留所绑定的 principal 正是这同一个 Arc(Arc::ptr_eq)的会话。同一用户在另一个入站上的会话持有另一个 principal,不受影响;
  3. 对其中每个会话执行策略(enforce):Keep 什么也不做,Close 取消会话的 token,CloseAfter(grace) 在 task tracker 上 spawn close_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 和各种策略的机制见用户与会话。

set_users、upsert_user 和 remove_user 不经过规划。它们发送 Command::Users,actor 运行 edit_users,即针对单个用户集的一套更小的准备和提交:

  1. 克隆运行中的 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),不做任何改变。
  2. 准备:对 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") 再次查找被编辑的用户集。
  3. 提交:commit_keys、publish_speed_limits,按各入站的策略撤销其被移除的 principal,然后存入每张用户表及其准入记录,最后把编辑后的 spec 记为运行中的 spec。

plane、监听器、handler 和版本都不受影响,也不生成 ApplyReport。remove_user 返回该用户是否在集合中;set_users 和 upsert_user 返回 ()。

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。在此版本中,仓库内的前端程序都没有调用它。

资源 保留条件 保留下来的内容
监听器 它的 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 只在步骤遍历中运行,位于所有切换和存储之后 无
位置 结果 之后的状态
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 句柄被丢弃时,它以零宽限期自行关闭。
  • 取消调用。 在 apply future 等待通道容量时丢弃它,什么也不会发送。命令一旦入队,无论调用方怎样,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 的回退路径,都没有专门的测试。修改这些路径时请补上。