Planning and applying a change
Source files: 41 · checked against 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
A supervisor never restarts to change what it runs. Every change, including the first spec, goes through one reconcile: plan compares what runs with the desired spec and lists the steps between them, the actor prepares everything that can fail without serving any of it, and then commits the rest, which cannot fail. A spec that is refused anywhere before the commit leaves the listeners, the data plane, the users and the speed limits exactly as they were. This replaces the etemenanki 2.x model, where a reload built a whole new generation of listeners, outbounds and routes: the reconcile works per resource, and a resource whose spec did not change keeps running, with its socket and its state.
This page is for contributors who change supervisor/src/topology/spec_plan/plan.rs, the apply path of supervisor/src/supervisor.rs, or supervisor/src/build/apply.rs. It covers how each resource is planned, outbound versions and drains, disruptive changes, the prepare and commit phases step by step, user revocation, ApplyReport, and update. What the crate owns overall and the actor’s other commands are on The supervisor; the spec types are on The spec; the rules plan checks first are on Validation.
Responsibilities
Section titled “Responsibilities”| Phase | Code | Can fail | Visible effect |
|---|---|---|---|
| Plan | topology/spec_plan/plan.rs → plan |
Yes: validate refuses the spec |
None. Pure: no sockets, no files, no clock. |
| Disruption gate | supervisor.rs → Actor::apply |
Yes: ApplyError::Disruptive |
None |
| Prepare | supervisor.rs → Actor::prepare |
Yes: ApplyError::Build, ApplyError::Bind |
Nothing is served yet. Sockets and devices opened here are not served until commit, and a refusal closes them. |
| Commit | supervisor.rs → Actor::commit |
No | Publishes the plane, starts and swaps listeners, stores user tables, stops what is gone, drains old outbounds, revokes removed users. |
The reconcile:
- decides, per resource, whether to keep it (
Reuse) or construct it (Build), and whether an inbound needs a new socket (Bind), a new handler on its old socket (SwapHandler), or nothing; - gives each build of an outbound or balancer tag its own version, so the old build can drain beside the new one;
- refuses a change that would end live connections unless the caller allowed it;
- publishes routes and outbounds together as one plane, only when DNS, an outbound, a balancer or the route changed;
- hands the live flows of a replaced or removed outbound version to its drain policy, and the live sessions of a removed user to the inbound’s removal policy;
- reports what it did in an
ApplyReport.
It does not:
- parse any format, read config files or decide when to reload. A front end lowers its own format into a
Specand callsapply; - define what a valid spec is.
plancallsbuild/validate.rs→validateand returns its error unchanged; - serve connections or route flows. It builds and swaps the objects that do; see Serving and The routing plane.
Key types
Section titled “Key types”RunningState and Plan
Section titled “RunningState and 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>;RunningStateis what the planner sees: the spec last applied (Nonebefore the first apply) and the latest version of every outbound and balancer tag ever applied.RunningState::empty()is{ spec: None, versions: empty }.versionsis never pruned. A removed tag keeps its entry, so a tag added back gets a version its old build, which may still be draining, never had.Plan::disruptions()yields(inbound, reason)for everyStep::Disrupt, in step order.planis public, like both types, but no front end at this revision calls it directly; they go throughSupervisor::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 },}| Step | Meaning | Work done in prepare | Work done in commit | Reported as |
|---|---|---|---|---|
Bind |
Open the listener an inbound asks for: a new inbound, or one whose BindSpec changed |
listener::bind, then Pending::pair |
Listener::start |
rebound, when the tag existed with another bind |
Build(r) |
Construct r from its spec: new, or changed |
Build the DNS resolvers, outbound, balancer, route table or inbound handler | Publish or swap it in | built |
Reuse(r) |
Keep the running r: its spec did not change |
Clone its Arc (outbounds, balancers, route, DNS) |
Nothing | reused |
PublishPlane |
Make the new route table and targets the ones new flows use, in one store | None | New Plane, epoch + 1, balancer probes |
Not reported |
SwapHandler |
Serve an inbound’s new connections with its new handler, on the listener it already has | Listener::prepare_swap |
Listener::swap |
swapped |
Disrupt |
The change ends the inbound’s live connections | None: the gate in Actor::apply refuses it first unless allowed |
None of its own | restarted |
StopAccepting |
Close the listener an inbound had on bind: the inbound was removed or moved |
None | Listener::stop and remove it |
Not reported |
CloseSessions |
Close the live sessions of an inbound tag that is gone | None | Sessions::close(Scope::Inbound) |
removed |
Drain |
Take an outbound version out of service | None | drain under the resolved policy |
drained |
Build(UserSet) and Reuse(UserSet) do no work of their own. Which user tables change is decided per inbound in prepare, by comparing admissions (see Prepare).
Resource and OutboundId
Section titled “Resource and 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 }| Value | 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 replaces identity by pointer address. Outbounds outlive a single apply, one tag can have an old version draining beside a new one, and a freed allocation can be reused by its successor, so anything that must tell two builds apart (a UDP fan-out link’s sub-links, a flow’s FlowEntry::outbound, the drain step) keys on the id, not on an Arc address. A balancer’s Resource carries only its tag; its versioned id is used for its Target.
ApplyOptions
Section titled “ApplyOptions”#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct ApplyOptions { pub allow_disruptive: bool,}allow_disruptive lets an apply execute a plan that contains Disrupt steps. It defaults to 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>,}| Field | Filled from | Meaning |
|---|---|---|
reused |
Step::Reuse |
Kept as they were: their spec did not change |
built |
Step::Build |
Constructed from their spec: new, or changed, in which case an outbound is a new version of its tag |
swapped |
Step::SwapHandler |
Inbounds that kept their listener and serve new connections with a new handler |
drained |
Step::Drain |
Outbound versions taken out of service; their live flows are handled by the drain policy |
rebound |
Step::Bind whose tag existed in the old spec with another bind |
Inbounds whose listener was closed and bound again |
restarted |
Step::Disrupt |
Inbounds whose change ended their live connections: the disruptive changes the apply was allowed to make |
removed |
Step::CloseSessions |
Inbound tags taken out of the spec, whose sessions were closed |
Every resource of the desired spec appears in reused or built, in plan order. Removed balancers and removed user sets appear nowhere; removed outbounds appear in drained. An inbound whose user table was replaced because its user set changed is reused, and the set is built.
ApplyError
Section titled “ApplyError”Every variant is defined in supervisor/src/build/apply.rs with thiserror. The rule variants come from validate; their rules are on Validation.
| Variant | Returned by | Display |
|---|---|---|
DuplicateTag { kind, tag } |
plan (validation) |
duplicate {kind} tag {tag} |
UnknownReference { from, kind, tag } |
plan (validation) |
{from} references unknown {kind} {tag} |
EmptyBalancer { balancer } |
plan (validation) |
balancer {balancer} has no members |
UnprobeableMember { balancer, member } |
plan (validation) |
balancer {balancer}: outbound {member} has no upstream a TCP health probe can reach |
Invalid { resource, reason } |
plan (validation); set_users, upsert_user, remove_user |
{resource}: {reason} |
Disruptive { inbound, reason } |
The gate in 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 |
Any call through the actor once it has gone | the supervisor has shut down |
Build and Bind carry the underlying io::Error as their #[source]. {bind} is the BindSpec Display: 127.0.0.1:1080 for TCP, udp 127.0.0.1:8443, unix:<path>, tun <name> (or tun auto) and tun fd <n>.
A Build error captured from the pinned etemenanki-app’s --test, which runs the same prepare (see check), for a Hysteria 2 inbound whose cert_file holds no PEM certificate:
configuration invalid: building inbound hy2-in failed: hysteria2: the certificate file contains no certificatesconfiguration invalid: is the prefix app/src/main.rs → main logs a failed --test with; the rest is ApplyError’s 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 is the default. The reason the code gives is that a TCP flow cannot migrate to the successor mid-stream, so by default the supervisor does not close the old version’s flows. UserRemovalPolicy (Keep, Close by default, CloseAfter(Duration)) is on Users and sessions. Each has a supervisor-wide default in Spec::policies, which an outbound (OutboundSpec::drain) or an inbound (InboundSpec::user_removal) may override.
What the actor keeps between applies
Section titled “What the actor keeps between applies”The reconcile’s state lives in the private Actor (supervisor/src/supervisor.rs), owned by the supervisor’s single task:
| Field | Type | Holds |
|---|---|---|
state |
RunningState<U> |
The spec last committed and every tag’s latest version |
shared |
Shared |
What every accept loop shares: the plane cell (PlaneCell), the Sessions registry, the flow Tracker and the TaskTracker that holds every accept loop and connection, and that balancer probes and grace timers are spawned on too |
root |
CancellationToken |
Parent of every listener, probe and session token |
dns, dns_target |
Option<Dns>, Option<Arc<Target>> |
The resolvers of the running DnsSpec, and the internal target of the DNS service when there is one |
targets |
HashMap<CompactString, Arc<Target>> |
The current version of each outbound and balancer tag |
probes |
HashMap<CompactString, CancellationToken> |
The probe token of each running balancer |
routes |
Option<Arc<CompiledRoutes>> |
The running route table |
listeners |
Vec<Listener> |
Every bound listener, found by BindSpec |
admissions |
HashMap<CompactString, Admissions<U>> |
Who each inbound admits, by inbound tag: the principal and credential per user |
keys |
UserKeys<U> |
Each admitted user id’s UserKey |
usage |
Arc<UsageBook<U>> |
The usage ledger and the ids its keys name; see Usage |
sampler |
Option<Sampler> |
The stats sampler, until SupervisorBuilder::start takes it to spawn it |
background, background_stop |
Vec<JoinHandle<()>>, CancellationToken |
The sampler and usage-sink tasks, and the token that stops them at the end of shutdown |
epoch |
u64 |
How many planes were published |
socket |
SocketOptions |
The socket policy every outbound socket is opened under; fixed for the supervisor’s lifetime |
A commit replaces state, dns, dns_target, targets, probes, routes, listeners, admissions, keys and epoch, stores into the plane cell of shared, and publishes new key ids into usage. sampler, background, background_stop and socket are not touched by an apply; they are described on The supervisor.
What prepare builds for one apply, not yet visible to anything:
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 carries each built balancer with its probe_interval and probe_timeout. swaps and users hold positions in Actor::listeners.
Data flow
Section titled “Data flow”One apply
Section titled “One apply”flowchart TB
caller["Supervisor::apply_with"]
chan["Command::Apply over the command channel"]
plan["plan: validate, then compare"]
gate{"Disrupt step without allow_disruptive?"}
prepare["Actor::prepare: build, bind, pair"]
commit["Actor::commit"]
ok["Ok: ApplyReport"]
err["Err: ApplyError, nothing changed"]
caller --> chan --> plan
plan -->|"a rule is broken"| err
plan --> gate
gate -->|"yes"| err
gate -->|"no"| prepare
prepare -->|"build or bind failed"| err
prepare --> commit --> ok
The entry points that reconcile all end in 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>;}Two more functions sit beside them and do not go through Actor::apply: epoch reads the plane cell, and check runs plan and prepare on an actor of its own.
impl<U: UserId> Supervisor<U> { pub fn epoch(&self) -> u64;}
pub async fn check<U: UserId>(spec: &Spec<U>) -> Result<(), ApplyError>;The handle reaches the actor through a private command channel. Every command carries a oneshot sender for its answer, and Supervisor::ask is the one helper that sends a command and awaits the answer:
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::startfirst asserts that the sample interval is not zero (the sample interval must not be zero). It then creates theActorand runs the first apply on the caller’s task, with the builder’soptions, before the actor task exists. A refused first spec returns itsApplyError: the actor task, the sampler and the usage sink are never started. Only after a successful first apply doesstartspawn the sampler and, when a sink is set, the usage push task, both underbackground_stop; then it creates the command channel (mpsc::channel(16)) and spawnsActor::run. A plan fromRunningState::empty()never contains aDisruptstep, sooptionshas no effect on the first apply at this revision.apply(spec)isapply_with(spec, ApplyOptions::default()). Both sendCommand::Apply(spec, options, reply)throughaskand await the reply.Actor::applyrunsplan(&self.state, &spec), then the gate, thenprepare(&plan, &spec, true), thencommit(plan, spec, prepared). The actor handles one command at a time (Actor::run), so two applies never interleave, and user edits, closes and session listings wait until the whole apply, prepare and commit, has finished. The serialisation comes from the one task that reads the channel, not from a lock.epoch()readsPlane::epochfrom the plane cell without going through the actor.Plane::empty()has epoch0, and eachPublishPlaneadds one. A handle exists only after a successful first apply, and the first plan always builds DNS and so always publishes, so a started supervisor reads1or more.
update and check are described below and under prepare.
Step order
Section titled “Step order”plan emits steps in one fixed order. Commit does not replay the list one step at a time: it publishes the plane, revokes users, starts listeners and swaps handlers in a fixed prologue (see Commit), then walks the list for the stops, closes and drains, which therefore run in the order shown.
flowchart LR dns["Dns"] --> out["each outbound"] --> bal["each balancer"] --> route["Route"] --> sets["each user set"] --> inb["each inbound: Bind, Build or Reuse"] inb --> pub["PublishPlane, if the plane changed"] --> swap["SwapHandler and Disrupt"] --> stop["StopAccepting"] --> close["CloseSessions"] --> drain["Drain"]
- Build and Reuse steps follow the desired spec’s order within each kind: outbounds in
Spec::outboundsorder, and so on. PublishPlanecomes after every build, so the plane is published with all of them.Drainsteps come last, afterPublishPlane: an old version is taken out only once the plane that replaces it is published.StopAcceptingprecedesCloseSessions, so a removed inbound’s listener is cancelled before its sessions are swept.Sessions::openchecks the listener’s stop token under the registry lock andSessions::closesweeps under the same lock, so a connection being accepted at that moment either is swept or never opens a session (see Users and sessions).
How plan decides
Section titled “How plan decides”plan first calls validate(desired) and returns its error. After that nothing can fail. It walks the resources in the order above, keeping plane_changed, the list of rebuilt outbound tags (rebuilt), the pending swaps and the pending drains.
dns_changed is true when there is no running spec or old.dns != desired.dns. It emits Build(Dns) or Reuse(Dns), and a changed DNS spec sets plane_changed, because the plane holds the DNS service’s target.
Outbounds
Section titled “Outbounds”For each desired outbound, prev is the running spec’s outbound with the same tag and running_version is running.versions[tag], looked up only when prev exists.
| Case | Steps | Version |
|---|---|---|
prev == Some(outbound) and not (dns_changed and the outbound uses DNS) |
Reuse(Outbound(tag@v)) |
Unchanged |
| Changed spec, or DNS changed and the outbound is not a blackhole | Build(Outbound(tag@v+1)), and Drain { tag@v, policy } queued |
next_version |
| New tag | Build(Outbound(tag@next)) |
next_version: 1 for a tag never seen, otherwise one more than the last version the tag had |
Every build sets plane_changed and adds the tag to rebuilt. Every outbound but a blackhole resolves names, so it holds the resolvers a changed DnsSpec replaces; a DNS change therefore rebuilds all of them (uses_dns = !matches!(protocol, Blackhole)).
An outbound in the running spec but not in the desired one gets Drain { tag@v, policy } and sets plane_changed. Its entry stays in versions.
drain_policy(tag) resolves the policy once, while planning. It takes the first outbound with that tag, looking in the desired spec first and then in the running one, and uses its OutboundSpec::drain override, or desired.policies.drain when that outbound has none. A changed outbound therefore drains its old version under the new spec’s override (or the new default), and a removed one under its old override.
Balancers
Section titled “Balancers”A balancer is reused when its BalancerSpec is unchanged, it has a running version, and none of its members is in rebuilt: a balancer holds its members’ targets, so a rebuilt member rebuilds it, and its health is learned afresh. Otherwise it gets Build(Balancer(tag)) with next_version(tag) and sets plane_changed. A balancer in the running spec but not the desired one sets plane_changed and keeps its version. Balancers have no Drain step (see Drains).
Outbound and balancer tags share one version map, as they share one namespace in a spec. A tag that turns from an outbound into a balancer, or back, continues the same count.
route_changed is true without a running spec or when old.route != desired.route. It emits Build(Route) or Reuse(Route) and a change sets plane_changed.
User sets
Section titled “User sets”Each desired set is Reuse(UserSet) when the running spec holds an equal set with its tag, and Build(UserSet) otherwise. User sets never set plane_changed and never touch an inbound’s steps. A removed set produces no step.
Inbounds
Section titled “Inbounds”Inbounds are matched by BindSpec, not by tag: the bind is the listener’s identity.
Running inbound on the same BindSpec |
Steps |
|---|---|
Equal InboundSpec |
Reuse(Inbound(tag)) |
Different InboundSpec |
Build(Inbound(tag)); queued: SwapHandler { tag }, then Disrupt when swap_disruption names a reason |
| None | Bind { tag, bind }, Build(Inbound(tag)); queued: a Disrupt when the change replaces a running TUN inbound’s device (see Disruptive changes) |
InboundSpec equality covers every field (tag, bind, sniff, protocol, users, user_removal), so renaming an inbound on the same bind, or changing its removal policy, is a handler swap. Inbound changes never set plane_changed.
Then, over the running inbounds:
- every running
BindSpecthat no desired inbound has getsStopAccepting { tag, bind }; - every running tag that no desired inbound has gets
CloseSessions { tag }.
So a moved inbound (same tag, new bind) is a Bind plus a StopAccepting and keeps its sessions, because sessions belong to the tag. A renamed inbound on the same bind is a SwapHandler plus a CloseSessions of the old tag, and keeps its listener. A removed inbound is a StopAccepting plus a CloseSessions.
Finally plan appends PublishPlane when plane_changed, then the queued swaps and disruptions, the stops, the closes and the drains, and returns the plan with the new versions.
Disruptive changes
Section titled “Disruptive changes”Two kinds of change end live connections by their nature, and plan marks them with Step::Disrupt:
| Change | Detected by | reason |
|---|---|---|
Any change to a TUN inbound that keeps its device (BindSpec), including its settings, sniff or tag |
swap_disruption: the running protocol is Tun |
the tun inbound's settings changed; its runtime restarts and its flows end |
A Hysteria 2 inbound on the same bind whose obfs changed |
swap_disruption: both Hysteria2 and old.obfs != new.obfs |
the hysteria2 obfuscation changed; connected clients cannot follow the new key |
| A TUN inbound’s device changed | plan, in the arm for a bind no running inbound has |
the tun device changed; the old device and its flows end |
Every change to a TUN inbound on the same device is a Disrupt: swap_disruption returns a reason whenever the running protocol is Tun. Any other change to a Hysteria 2 inbound on the same bind is a SwapHandler, not a Disrupt.
Actor::apply checks plan.disruptions().next() before preparing anything. Without allow_disruptive it returns ApplyError::Disruptive for the first disruption in step order, for example:
inbound hy2-in: the hysteria2 obfuscation changed; connected clients cannot follow the new key; this ends its live connections and needs allow_disruptiveHow the front ends choose:
- etemenanki-app calls
Supervisor::applyand so never allows a disruption.app/src/instance.rs→Core::reload_withlogsreload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply itand keeps what runs. A TUN change or a Hysteria 2 obfuscation change therefore needs a restart. The REST API and the FFI reload through the sameCore(etemenanki-app: running, reloading and shutting down). - katana’s
src/manager/node.rs→reconcileapplies each new node spec withApplyOptions { allow_disruptive: true }, so a Hysteria 2 obfuscation change from the panel is applied at once. See The node manager.
Worked examples
Section titled “Worked examples”Each row is a unit test in supervisor/tests/unit/plan.rs, starting from the spec it names:
| Change | Plan |
|---|---|
First apply of a SOCKS inbound admitting set people, outbounds direct and proxy, balancer lb over proxy |
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; versions direct 1, proxy 1, lb 1 |
| The same spec again | Seven Reuse steps, nothing else |
direct’s address family changed |
Build(direct@v2) … PublishPlane … Drain { direct@v1, Keep }; route and inbound reused |
| A route rule added | Build(Route), PublishPlane; DNS, outbounds and inbound reused; no drain |
A user added to people |
Build(UserSet people), Reuse(Inbound socks); no bind, stop, swap or publish |
| The SOCKS inbound turned into HTTP on the same bind | Build(Inbound socks), SwapHandler socks; no bind, stop, close or disruption |
| The SOCKS inbound moved to port 1081 | Bind socks 127.0.0.1:1081 before Build(Inbound socks), StopAccepting socks 127.0.0.1:1080; no swap, no close |
The DNS spec changed, with outbounds direct, hole (a blackhole) and proxy |
Build(Dns), Build(direct@v2), Reuse(hole@v1), Build(proxy@v2), PublishPlane, drains of direct@v1 and proxy@v1 |
Versions
Section titled “Versions”stateDiagram-v2 [*] --> Current: Build, published with the plane Current --> Current: Reuse keeps the same Arc Current --> Draining: Drain after a rebuild or removal Draining --> [*]: Keep, drain does nothing Draining --> [*]: Close, closed token cancelled at once Draining --> [*]: CloseAfter, closed token cancelled after the grace
- A tag’s versions count up from 1 and are never reused.
next_version(tag)isrunning.versions[tag] + 1, or1for a tag never applied. A removed tag keeps its version, so re-adding it buildsv+1. - A reused outbound keeps its version and the same
Arc<Target>, and with it the state its connector holds, such as a Hysteria 2 QUIC connection or a WireGuard tunnel. Building an outbound only constructs its connector; nothing is dialled until the first flow. - The actor stops holding an old version’s
Arc<Target>at the commit that replaces or removes it:self.targetsis replaced in commit step 3, and the old map is dropped when the commit returns. ACloseAfterdrain timer holds theArcuntil it fires or shutdown ends it. Each flow keeps its own link plus the version’sclosedtoken, whichCloseorCloseAftercancels (see Drains). - The supervisor makes two targets of its own, both with the empty tag, which no spec tag can be (
validaterefuses an empty tag):Plane::empty()’s blackhole,@v0, which drops every flow before the first apply, and the DNS service’s target, whose version isepoch + 1when it is built. A DNS rebuild always publishes a plane, so that is the epoch of the plane that first carries it.
Prepare
Section titled “Prepare”async fn prepare(&mut self, plan: &Plan, spec: &Spec<U>, bind: bool) -> Result<Prepared<U>, ApplyError>;Actor::prepare does everything that can fail. It is async because it awaits listener::bind, which opens a TUN device asynchronously. It builds into locals and a staged copy of the user keys, and builds the Prepared only at its end, so returning an error at any step drops everything built so far and leaves the actor unchanged.
| # | Step | What it does | Error |
|---|---|---|---|
| 1 | DNS | With Build(Dns) in the plan, or no running DNS: build_dns(&spec.dns, plane, &socket), where plane is a Weak reference to the plane cell (the through-proxy half of a split resolver dials through it), and a new internal target Target::outbound(internal_id(epoch + 1), Outbound::Dns(service)) when the spec has a DNS service. Otherwise clone the running Dns and target. |
Build { resource: Dns } |
| 2 | Outbounds, in plan order | Build(Outbound(id)): build_outbound(outbound, &dns, &socket), wrapped as Target::outbound(id, built). Reuse(Outbound(id)): clone self.targets[tag]. |
Build { resource: Outbound(..) } |
| 3 | Balancers, in plan order | Build(Balancer(tag)): one Member::new(member, targets[member], probe_target(outbound)) per member, over the targets from step 2, then Balancer::new(members, strategy), wrapped as Target::balancer(tag@plan.versions[tag], ..) and queued in new_balancers. Reuse: clone the running target. |
Build { resource: Balancer(tag) }: a balancer needs at least one outbound, which validation already rules out |
| 4 | Route | With Build(Route) or no running table: compile_routes(&spec.route) (see The routing plane). Otherwise clone the running Arc. |
Build { resource: Route } |
| 5 | Admissions, per desired inbound | admit(credential_kind(protocol), set, self.admissions.get(tag), &mut keys) on a clone of self.keys; a user admitted without a principal they keep gets a new one, Principal::user(keys.key(id), id.name()). Principals no longer admitted are queued in revoked with removal_policy(spec, inbound). |
None |
| 6 | Inbounds, per desired inbound | See below. | Build { resource: Inbound(tag) } or Bind |
The member lookups in step 3 use expect("validated: members are outbounds") and expect("validated: members are probeable"); the spec lookups in steps 2 and 3 use expect("the plan builds the spec's outbounds") and expect("the plan builds the spec's balancers"). Plain indexing, which panics on a missing key, relies on the same guarantees: self.targets[&id.tag] for a reused outbound, targets[member] and plan.versions[tag] for a built balancer, self.targets[tag] for a reused one, and admissions[&inbound.tag] in step 6. Validation and plan guarantee all of them.
Step 6, for each desired inbound, with running the position of the listener whose bind equals the inbound’s:
Build(Inbound(tag))in the plan. The handler is built with the principals from step 5:build_handler(inbound, admitted, circuits).circuitsis the running Hysteria 2 listener’s circuit budget (Listener::circuits()) when both the running and the desired inbound on that bind are Hysteria 2 with the samemax_circuits(same_circuits), so live circuits keep counting against it; otherwise the handler gets a newSemaphoreof the new size. Then:- with a running listener,
Listener::prepare_swap(handler)refuses a handler of another kind, checks the Hysteria 2 handler’s Salamander key (check_obfs, which turns aSalamander::newerror into anInvalidInputerror with the same text), and for TUN duplicates the listener’s device (try_clone) and prepares the duplicate for the runtime the swap restarts (TunInbound::prepare). The result is queued inswapswith the listener’s position; - without one, in a dry run (
bind == false), nothing more; - otherwise
listener::bind(tag, bind).awaitopens the socket or device, andPending::pair(bound, handler)pairs it with the handler and builds what serving needs. The pair is queued instarted.
- with a running listener,
- Otherwise (the inbound is reused), when a listener runs on the bind and the admissions differ from the running ones (none recorded, or
same_admissionsis false), the user table alone is rebuilt:user_table(inbound, admitted), thenListener::prepare_users(table), which refuses a table of another kind. The result is queued inusers.
What listener::bind and Pending::pair do per kind (the details are on Serving):
| Bind | listener::bind |
Pending::pair |
|---|---|---|
| TCP | std::net::TcpListener::bind, non-blocking, wrapped for tokio |
Pairs the listener with the StreamInbound |
| Unix | Removes a socket file already at the path, taken to be stale from a crashed run; refuses anything else there with <path> exists and is not a socket (AlreadyExists); binds and records the file’s device and inode in a SocketFile |
As TCP |
| UDP (Hysteria 2) | std::net::UdpSocket::bind |
check_obfs, then Hy2Inbound::new and Hy2Inbound::open(socket), which stores the Salamander key in the new inbound and opens the QUIC endpoint on the socket. The endpoint reads nothing off the socket until it is run. |
| TUN | tun::open for a created device, tun::adopt for a supplied descriptor, then a first duplicate with try_clone |
TunInbound::prepare(first), which registers the duplicate with the tokio reactor (AsyncFd) and configures its IP stack. Nothing reads the device until it is run. |
Any other pairing of bind and handler is refused with the handler is not of the kind its listener serves.
same_admissions holds when both maps have the same users in the same order with the same principal (Arc::ptr_eq) and the same credential, so the running table already says what a new one would. A user whose only change is their speed limit keeps both.
listener::bind’s error becomes ApplyError::Bind { inbound, bind, source }. Every other failure in step 6 is ApplyError::Build { resource: Inbound(tag), source }, for example:
- a certificate or key that does not parse, such as
hysteria2: the certificate file contains no certificates(the example above); inbound <tag>: protocol and bind do not match, frombuild_handler;the handler is not of the kind its listener serves, fromprepare_swap,prepare_usersorPending::pair;- a Salamander key
check_obfsrefuses, or a failure ofHy2Inbound::openorTunInbound::prepare.
build_handler also calls Masquerade::new for a Hysteria 2 masquerade, but validation calls it first, so a masquerade it refuses never reaches prepare.
What an early return from prepare releases (and what check drops with the Prepared it gets):
- a TCP or UDP socket bound in step 6 is closed. A Unix listener’s socket file is removed as well (
SocketFile’sDrop, only while the path still holds the file this listener created; a failed removal logscould not remove <path>: <e>atdebug). A created TUN device’s descriptors are closed; - a Hysteria 2 endpoint opened by
Pending::pairis dropped with its newHy2Inbound, which ends the endpoint’s driver task. A TUNPreparedDevice, fromPending::pairor fromprepare_swap, is dropped and its duplicate descriptor closed; - new outbounds, balancers, resolvers and route tables are freed. None of them has started a task or opened a connection;
- the staged
UserKeyscopy is dropped, so no key handed out in step 5 is published.
The disruption gate runs before prepare, so a spec it refuses has had nothing built.
check: a dry run
Section titled “check: a dry run”check(spec) runs plan against RunningState::empty() in a fresh Actor (default socket policy), then prepare(&plan, spec, false), and drops the result. Its doc comment says it “refuses what a start would, short of a port or device that cannot be bound”. It builds everything a start would: resolvers, outbounds, balancers, the route table, every handler with its certificates, every user table. It binds nothing and does not pair, so a port or device that cannot be bound, and the checks Pending::pair makes (check_obfs, Hy2Inbound::open, TunInbound::prepare), are not reached. Because it plans from nothing, it never meets the disruption gate. etemenanki-app’s --test (app/src/instance.rs → check) and katana’s per-node check (src/runtime.rs → check_node) both call it; see Validation.
Commit
Section titled “Commit”Actor::commit(&mut self, plan: Plan, spec: Spec<U>, prepared: Prepared<U>) -> ApplyReport cannot fail. Whatever could fail was done while binding, building and pairing; the remaining expect and unreachable! calls state what prepare guaranteed (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), and so does the plain index self.targets[tag] in the slot fill of step 4.
| # | Step | Why here |
|---|---|---|
| 1 | commit_keys(keys): adopt the staged keys and publish the ids of new keys to the usage book (UsageBook::names) |
Every key must be published before a session can bind it: the usage book panics on a key it cannot name (… was bound before it was published) |
| 2 | publish_speed_limits(&spec): every user’s speed_limit to Tracker::set_speed_limits; a user in several sets gets the smallest; a user without a UserKey (self.keys.get(id) is None) is skipped; set_speed_limits sets the rate of every pacer it has, unlimited for a key not in the map, and inserts a pacer for each limited key that has none |
Only at commit, so a refused apply never changes live limits. Live flows feel a new limit on their next read or write. |
| 3 | Replace self.targets with the prepared map, keeping the old one as old_targets |
The drains in step 11 need the old versions |
| 4 | With PublishPlane: epoch += 1; build Plane::new(routes, slots, dns_target, epoch), where slots holds one target per route slot, in CompiledRoutes::targets() order (self.targets[tag] for each; Plane::new has a debug_assert_eq! that the two lengths match); store it in the plane cell |
One store swaps routes, outbounds and the DNS service together, so a rule never names a tag the plane lacks |
| 5 | With PublishPlane: cancel the probe token of every balancer that is not Reuse in the plan; for each new balancer, a child of the root token and Balancer::spawn_probe on the task tracker, with the new dns.servers resolver and TcpDialer::new(socket) |
A balancer no longer in the plane stops probing; a new one probes against the resolver its members dial with |
| 6 | Store routes, dns and dns_target |
Kept for the next apply’s Reuse |
| 7 | Sessions::revoke(principals, policy) for every revoked group |
Before any new user table is stored (see User revocation) |
| 8 | Listener::start for every started pair, each with a child of the root token, appended to listeners |
New listeners route through the new plane from their first connection |
| 9 | Listener::swap for every swaps entry |
Positions recorded in prepare are still valid: step 8 only appends |
| 10 | Listener::store_users for every users entry |
After the revocations |
| 11 | Walk the plan’s steps in order, filling the report; StopAccepting stops and removes the listener on its bind (swap_remove); CloseSessions closes the tag’s sessions; Drain hands the old target to drain when old_targets[tag] is exactly that version, and pushes the step’s id to report.drained either way |
After every swap and table store, so swap_remove cannot move a position still to be used; drains after the publish |
| 12 | self.admissions = admissions; self.state = RunningState { spec, versions: plan.versions } |
The next plan starts from here |
| 13 | tracing::debug!("applied: {} built, {} reused, {} swapped, {} drained", …) |
Log target etemenanki_supervisor::supervisor |
Without PublishPlane (a user-set or inbound-only change) the plane, the epoch and the probes are left alone. Plane::empty(), the plane before the first commit, routes everything through CompiledRoutes::single over the empty tag to its one slot, the @v0 blackhole.
A plane holds only the targets its route table names. An outbound that only a balancer names is reached through that balancer’s Target, which holds the member targets built or reused in prepare.
What the listener operations do
Section titled “What the listener operations do”These live in supervisor/src/system/listener.rs; the accept loops they drive are on Serving.
| Operation | Stream listener (TCP, Unix) | Hysteria 2 | TUN |
|---|---|---|---|
Listener::start |
Spawns run_stream_inbound on the task tracker with a watch channel of the handler |
Spawns run_hysteria_inbound with the endpoint and an ArcSwap of the tag |
Spawns run_tun_inbound with an unbounded restart channel |
Listener::swap |
Sets the listener’s tag to the new inbound’s, then send_replace of the new Arc<StreamInbound> into the watch channel; connections accepted after the swap are served by it |
set_obfs when the key changed; sets the listener’s tag and stores the new tag in the ArcSwap that run_hysteria_inbound reads once per connection, so a renamed inbound’s new connections are filed under the new tag; replaces the kept circuit budget with the handler’s; then set_connection_config, set_quic and set_max_connections on the running Hy2Inbound |
Sets the listener’s tag, then sends the new handler and the prepared duplicate on the restart channel; if the send fails, spawn_tun starts a runtime task with them and its sender replaces the old one |
Listener::store_users |
Stores the new user table (UserTable::store_into); see Users and sessions |
set_authenticator |
Nothing: a TUN device admits no users |
Listener::stop |
Cancels the listener’s token, so the accept loop stops. Sessions run under children of the root token, not of this one. | Cancels the listener’s token; how the endpoint winds down is on Serving | Cancels the listener’s token, of which each run of the runtime task holds a child |
Every started listener logs inbound <tag> listening on <bind> at info (target etemenanki_supervisor::system::listener). A created TUN device also logs inbound <tag> owns tun device <name>, and a supplied one inbound <tag> serves a supplied tun device, when it is bound in prepare.
Drains
Section titled “Drains”A Drain step resolves to one call:
fn drain(target: Arc<Target>, policy: DrainPolicy, root: &CancellationToken, tracker: &TaskTracker);| Policy | What drain does to the old version’s flows |
|---|---|
Keep (default) |
Nothing: the supervisor does not close them |
Close |
Target::close_flows(): closes them at once |
CloseAfter(grace) |
Spawns a task on the tracker that selects on root.cancelled() and sleep(grace); after the sleep it calls close_flows(), so it closes them after grace |
New flows follow the published plane whatever the policy: to the tag’s new version after a rebuild, or wherever the new route table sends them after a removal.
How closing works (supervisor/src/topology/plane.rs):
- Every
Targetowns aclosed: CancellationToken.close_flows()cancels it. Target::connect_streamandTarget::connect_datagramresolve the target first (a balancer picks a member) and wrap the opened link inGuarded<T>with the resolved target’sclosedtoken. A resolved target that is itself a balancer fails the call withbalancer member <id> is itself a balancer, which validation rules out.Guardedchecks the token before everypoll_read,poll_write,poll_flush,poll_send_toandpoll_recv_from, including a receive that is already waiting, and returnsio::ErrorKind::ConnectionAbortedwith the textthe outbound this flow was opened on was drained.poll_shutdownis passed through unchecked. The runtime relaying the flow sees an outbound error and ends that flow.- A flow through a balancer is opened on the member chosen for it, so it carries that member’s token and is closed by that member’s drain. Balancers therefore have no
Drainstep: a rebuilt balancer whose members were reused does not close their flows, and a rebuilt member drains its own. - The DNS service’s internal target has no
Drainstep either, so the supervisor does not close flows already handed to the old service. - A pending
CloseAfterdrain timer does not holdshutdownpast its owngrace: the timer runs on the task tracker, whichshutdownwaits on for at most itsgrace, and it selects on the root token, whichshutdownthen cancels (shutdown_ends_pending_grace_timers).
Drain policies are resolved while planning (see Outbounds) and fixed in the step, so a later apply cannot change the policy of a version already drained.
User revocation
Section titled “User revocation”Users change through an apply (a changed UserSet, an inbound that now admits another set or reads another credential kind) or through the user path below. Either way build/users.rs → admit compares an inbound’s running admissions with the new ones:
- a user who is still admitted, and still holds the exact credential they were admitted by before, keeps their
Arc<Principal>, even when the inbound now reads another credential kind, so their live sessions are untouched; - a user who is no longer admitted, or whose earlier credential changed or is gone, loses their old principal, which is returned in
revoked. If they are still admitted, they get a new principal,Principal::user(keys.key(id), id.name()).
UserKeys::key returns the key a user id already has, or hands out the next one. Keys count up from 1, so key n names ids[n - 1] (build/users.rs → id_of).
Each inbound’s revoked principals are queued with that inbound’s policy: inbound.user_removal, or spec.policies.user_removal when unset. Admissions are keyed by inbound tag, so a renamed inbound starts with fresh principals; its old tag’s sessions are closed by CloseSessions rather than revoked.
At commit, revocation runs before any new user table is stored, and before new listeners start or handlers swap:
sequenceDiagram participant C as Actor commit participant P as old Principal participant H as handshake on the old table participant S as its Session C->>P: Sessions::revoke marks it revoked C->>S: live sessions bound to it follow the policy H->>P: authenticates as the old principal H->>S: first flow calls Session::admit S-->>H: refused, the principal was revoked C->>C: store_users installs the new table
Sessions::revoke(principals, policy) takes the registry lock and, for each principal:
- marks it revoked under the policy (
Principal::revoke; only the first revocation counts); - looks up the sessions of its
UserKeyin the registry’sby_userindex, and keeps those whose bound principal is this sameArc(Arc::ptr_eq). The same user’s sessions on another inbound hold another principal and are left alone; - applies the policy to each of them (
enforce):Keepdoes nothing,Closecancels the session’s token, andCloseAfter(grace)spawnsclose_afteron the task tracker, which selects on the session’s token andsleep(grace)and cancels the token after the sleep.
A session whose first flow presents a principal that is already revoked is refused: Session::admit checks the principal under the same registry lock when the first flow binds the session, and refuses it under every policy with the user was removed before the session opened a flow, because a removal policy only spares sessions that were bound to the principal at the removal. Revoking before store_users makes this cover a handshake that authenticated against the table about to be replaced. The mechanics of Principal, Session and the policies are on Users and sessions.
The user path
Section titled “The user path”set_users, upsert_user and remove_user do not plan. They send Command::Users and the actor runs edit_users, a smaller prepare and commit over one user set:
- Clone the running spec; without one, answer
ApplyError::Stopped. Edit the named set:Setreplaces its users,Upsertinserts or replaces one,Removeremoves one. A missing set isApplyError::Invalid { resource: UserSet(tag), reason: "no such user set" }(user set <tag>: no such user set).SetandUpsertalways count as a change, even when the result is identical. ARemoveof a user who is not in the set returnsOk(false)and changes nothing. - Prepare, for every inbound whose
usersnames the set:validate_admission(inbound, set), which returnsOkat once for a TUN inbound (it has no credential kind);admiton a staged copy of the keys; the listener on the inbound’s bind (expect("every running inbound has its listener"));user_table;Listener::prepare_users. Every such inbound’s table is rebuilt, with nosame_admissionscheck. Any failure returns before anything is stored:validate_admission’s error as is, and an error fromuser_tableorprepare_usersasApplyError::Build { resource: Inbound(tag) }. The edited set is looked up again for this step withexpect("found above"). - Commit:
commit_keys,publish_speed_limits, revoke every inbound’s removed principals under its policy, then store each table and its admissions, then record the edited spec as the running spec.
The plane, the listeners, the handlers and the versions are not touched, and no ApplyReport is made. remove_user returns whether the user was in the set; set_users and upsert_user return ().
update and update_with
Section titled “update and update_with”update(edit) is update_with(edit, ApplyOptions::default()). It sends Command::Update(Box::new(edit), options, reply), where the boxed closure is an Edit<U> (see One apply). The actor clones its running spec, calls edit on the clone inside the actor task, and applies the result exactly as apply_with would. Because the edit runs on the actor, it always sees the spec the previous command left, with no read-modify-write race between two callers. It answers ApplyError::Stopped if there is no running spec, which a started supervisor always has. None of the in-tree front ends calls it at this revision.
What an apply keeps
Section titled “What an apply keeps”| Resource | Kept when | What survives |
|---|---|---|
| A listener | Its BindSpec is unchanged (Reuse or SwapHandler) |
The socket or device and its accept loop; for a stream listener the loop’s MAX_HANDSHAKES_PER_INBOUND (204,800) and MAX_LIVE_CONNECTIONS_PER_INBOUND (4,194,304) semaphores, which run_stream_inbound creates once; for Hysteria 2 the QUIC endpoint and, with the same max_circuits, the circuit budget |
| An outbound | Its spec is unchanged, and so is the DNS spec unless it is a blackhole | The same Arc, with the state its connector holds, such as a QUIC connection or a tunnel |
| DNS resolvers | DnsSpec is unchanged |
The resolvers and their shared answer cache, and the DNS service’s target |
| A balancer | Its spec and every member are unchanged | Its member health and its running probes |
| The route table | RouteSpec is unchanged |
The compiled table |
| A user’s principal | The user is still admitted by the same credential | Their live sessions; a speed-limit change keeps them too |
| The plane | Nothing in DNS, outbounds, balancers or the route changed | The published Plane and the epoch |
Sessions run under children of the root token, not of their listener’s token, so stopping a listener or swapping its handler does not cancel them. The reconcile closes sessions through CloseSessions, for an inbound tag that left the spec, and through Sessions::revoke, which applies a removal policy to the sessions bound to a revoked principal. A Disrupt step marks a change whose reason says it ends the inbound’s live connections (see Disruptive changes).
Policies (Spec::policies) have no step. New defaults apply to the revocations and drains of the apply that brings them, and after.
A route change does not move an established flow: the connector loads the plane once per connect (supervisor/src/connector.rs → AppConnector::connect) and a UDP fan-out link once per packet it sends (FanOutLink::poll_send_to), so a new connection, each new mux sub-flow and each UDP packet of a fan-out link follow the plane current when they are routed (see The routing plane).
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| Planning has no side effects | plan takes references and does no I/O |
Every test in supervisor/tests/unit/plan.rs builds specs by hand and binds nothing |
| A spec refused anywhere before commit changes nothing | Validation in plan; the gate before prepare; prepare builds into locals and staged keys; speed limits and keys are published only at commit |
a_refused_spec_changes_nothing, an_invalid_spec_is_refused_without_a_plan |
| A listener bound for a refused spec is released | Prepare’s locals own the bound handle; an early return drops it | a_refused_spec_changes_nothing |
| A version of a tag is never reused | versions keeps removed tags; next_version counts up |
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 |
| An old version is drained only after the plane that replaces it is published | Drain steps after PublishPlane |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
| Only the version a step names is drained | target.id() == outbound check in commit; at this revision old_targets[tag] always holds the version the step names, so the check is defensive |
none |
| An unchanged outbound keeps its state | Reuse clones the running Arc |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
| A DNS change reaches everything that holds a resolver | uses_dns rebuild of every non-blackhole outbound |
a_dns_change_rebuilds_every_outbound_but_a_blackhole |
| A balancer is rebuilt with its members, and kept otherwise | members_kept |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is |
| Removing a balancer stops routes naming it | plane_changed on removal |
removing_a_balancer_alone_republishes_the_plane |
| User changes never touch the plane or the listeners | User sets never set plane_changed; tables are swapped in place |
a_user_set_change_alone_touches_neither_the_plane_nor_the_inbounds, a_speed_limit_change_keeps_the_users_sessions |
| An unchanged bind keeps its socket | Inbounds matched by 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 |
| A moved inbound keeps its sessions | Sessions belong to the tag | a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions |
| A removed or renamed inbound’s sessions are closed | CloseSessions per vanished tag, after 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 |
| Disruptive changes need consent | The gate in Actor::apply |
a_hysteria2_obfuscation_change_needs_allow_disruptive, a_tun_device_change_needs_allow_disruptive |
| A session whose first flow presents a revoked principal is refused | Session::admit checks the principal under the registry lock; revocation runs before 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 |
| A key is published before a session can bind it | commit_keys first in commit and in edit_users |
none; a violation panics in the usage book |
| The epoch moves only when a plane is published | epoch += 1 only under PublishPlane |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow (epoch + 1), a_refused_spec_changes_nothing (unchanged) |
| Grace timers end with the supervisor | Drain timers select on the root token; removal timers select on the session’s token, a child of the root | shutdown_ends_pending_grace_timers, cancelling_the_root_cancels_every_session |
| Positions recorded in prepare hold at commit | New listeners are appended; swap_remove runs only in the step walk, after every swap and store |
none |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Where | Result | State afterwards |
|---|---|---|
plan |
A validation ApplyError |
Unchanged |
| The gate | ApplyError::Disruptive naming the first disruption |
Unchanged; nothing was built |
| Prepare | ApplyError::Build or ApplyError::Bind for the first failure, in the order of the prepare table |
Unchanged; what was built is dropped and released |
| Commit | Cannot fail | The new spec runs |
| The actor is gone | ApplyError::Stopped from apply, update and the user methods; 0 from close; an empty list from sessions; shutdown returns |
The supervisor has shut down |
- Stopped.
Supervisor::askmaps both a closed command channel and a dropped reply toApplyError::Stopped. The actor returns afterCommand::Shutdown, and it shuts itself down with a zero grace when the lastSupervisorhandle is dropped. - Cancelling a call. Dropping an
applyfuture while it waits for channel capacity sends nothing. Once the command is queued, the actor runs the apply to the end whatever happens to the caller; the reply is then discarded. - Starting. A refused first spec fails
SupervisorBuilder::startwith theApplyError; the actor task, the sampler and the usage sink are never started, and prepare’s locals are dropped as above. A zero sample interval panics instartbefore the first apply. - No running spec.
Command::Updateandedit_usersanswerApplyError::Stoppedwhen the actor has no running spec. A started supervisor always has one, so this is a guard, not a path a caller can reach.
How the front ends wrap the result:
- etemenanki-app’s
Core::reload_withlogsreload: <error>; keeping the running configfor most refusals and the dedicated line above forDisruptive, and logsconfig reloaded: <summary>on success. The summary lists the report’s non-empty groups in the order built, swapped, drained, rebound, restarted, removed, reused, aslabel [item, item]joined by;, ornothing to runwhen every group is empty. See etemenanki-app: running, reloading and shutting down. - katana logs an apply’s error for the node; see The node manager.
Limits
Section titled “Limits”| Item | Value | Defined in |
|---|---|---|
| Command channel | mpsc::channel(16) |
SupervisorBuilder::start |
| Applies in flight | One: the actor serialises every command | Actor::run |
| Plane epoch | u64, 0 for Plane::empty(), +1 per publish |
Actor::commit |
| Outbound and balancer versions | u64 per tag, starting at 1 |
plan → next_version |
versions map |
One entry per outbound or balancer tag ever applied; never pruned | RunningState, Plan |
| Internal targets | Empty tag; @v0 for the empty plane, @v<epoch + 1> for the DNS service |
Plane::empty, internal_id |
| TUN restart channel | mpsc::unbounded_channel, written only by the actor, once per swap |
system/listener.rs → spawn_tun |
| Stream handler channel | watch::channel, one value |
Listener::start |
| Resource lookups | Linear scans by tag or by bind over the spec’s vectors | plan, Actor::prepare, Actor::commit |
supervisor/tests/unit/plan.rs drives plan alone. Its helpers build a SOCKS inbound on 127.0.0.1:1080 with a direct freedom outbound (spec()), and a full_spec() that adds a SOCKS client outbound proxy to 192.0.2.1, a failover balancer lb over it (probe every 30 s, timeout 5 s) and a user set people admitted by the inbound. A TUN inbound is only a DeviceSpec value, so nothing is opened.
| Test | Behaviour it pins |
|---|---|
the_first_plan_binds_and_builds_everything_then_publishes |
The exact step list and versions of a first apply |
an_unchanged_spec_reuses_everything_and_publishes_nothing |
Seven Reuse steps, one per resource; versions unchanged |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
v2 is built, v1 drained under Keep, after PublishPlane; route and inbound reused |
a_drain_takes_the_supervisor_default_unless_the_outbound_overrides_it |
policies.drain = Close applies; an outbound’s CloseAfter(5 s) override wins |
each_rebuild_of_a_tag_takes_the_next_version |
Two rebuilds give v3 and drain v2 |
a_removed_outbound_drains_under_its_own_policy |
Its old Close override applies; the plane is republished; its version is kept |
a_tag_added_back_takes_a_version_it_never_had |
Remove and re-add gives v2 |
a_dns_change_rebuilds_every_outbound_but_a_blackhole |
Freedom and proxy rebuilt and drained; the blackhole reused |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is |
A non-member rebuild keeps lb@v1; a member rebuild builds lb@v2 and publishes |
removing_a_balancer_alone_republishes_the_plane |
Publish, no drain, version kept |
a_route_change_alone_rebuilds_only_the_route_table |
Only Build(Route) and a publish |
a_user_set_change_alone_touches_neither_the_plane_nor_the_inbounds |
Build(UserSet), Reuse(Inbound), no publish, bind, stop or swap |
a_protocol_change_on_the_same_bind_swaps_the_handler_and_keeps_the_listener |
SwapHandler, no bind, stop, close or disruption |
a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions |
Bind before Build, StopAccepting of the old bind, no close |
a_removed_inbound_stops_accepting_and_closes_its_sessions |
The plan ends with StopAccepting then CloseSessions |
an_inbound_renamed_on_the_same_bind_swaps_and_closes_the_old_tags_sessions |
SwapHandler for the new tag, CloseSessions for the old, no bind or stop |
a_tun_device_change_opens_a_new_device_and_disrupts |
A changed DeviceSpec is a Bind, a StopAccepting and a Disrupt, with no swap |
a_tun_settings_change_on_the_same_device_swaps_and_disrupts |
A changed TunSpec is a SwapHandler and a Disrupt |
a_hysteria2_obfs_change_swaps_and_disrupts |
A new Salamander key is a SwapHandler and a Disrupt |
a_hysteria2_masquerade_change_swaps_without_disrupting |
A masquerade change is a SwapHandler with no Disrupt |
an_invalid_spec_is_refused_without_a_plan |
An unknown default outbound returns UnknownReference { from: Route } |
supervisor/tests/hot_swap.rs runs real supervisors on loopback sockets with the helpers in supervisor/tests/support/mod.rs (base_spec(): outbounds direct and block, default route direct, system DNS; SOON = 5 s, QUIET = 400 ms):
| Test | Behaviour it pins |
|---|---|
an_established_connection_survives_an_apply_that_keeps_its_inbound |
An apply adding an outbound and a rule reuses the inbound; the open SOCKS connection keeps echoing |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
The epoch grows by one; the next packet of a live SOCKS UDP association takes the new route; the open TCP flow does not; a new TCP flow does |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
On one live VLESS mux session, a sub-flow opened after the change takes the new route and the earlier one keeps its own |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close |
With the freedom outbound direct: under Keep, a flow opened on direct@v1 still echoes after the apply that builds v2; under Close, the flow on v2 is closed, the one on v1 still echoes (v1 was not drained again), and a new flow uses v3 |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
A reused Hysteria 2 outbound opens no second QUIC connection; a rebuilt one dials anew while the flow on the old one keeps going |
a_hysteria2_obfuscation_change_needs_allow_disruptive |
Refused with Disruptive; allowed, it reports restarted and no rebound |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
SOCKS to HTTP on one port: swapped, no rebound; the open SOCKS connection keeps relaying; a new one speaks HTTP |
a_refused_spec_changes_nothing |
After an UnknownReference and a Bind refusal: same epoch, the old route and users, and the listener bound while preparing is released |
removing_a_user_closes_only_their_sessions_and_refuses_them_after |
Under Close and under Keep: only the removed user’s session is closed (or none), their next SOCKS handshake fails, and close(Selector::User) closes the other user’s |
a_speed_limit_change_keeps_the_users_sessions |
set_users with a new limit keeps the session and its id |
removing_an_inbound_closes_its_sessions |
removed lists the tag; its connection closes, the other inbound’s does not; its port refuses connections |
a_tun_device_change_needs_allow_disruptive |
A device change and a settings change are both refused without the option. Needs CAP_NET_ADMIN; skipped when a device cannot be created |
shutdown_ends_pending_grace_timers |
An hour-long removal grace and an hour-long drain grace do not hold up shutdown(Duration::ZERO) |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
shutdown(Duration::ZERO) returns within SOON (5 s) while a client is stuck in its QUIC handshake; it does not wait for the handshake’s idle timeout (see The supervisor) |
Other tests on this path:
| File | Test | Behaviour it pins |
|---|---|---|
supervisor/tests/unit/plane.rs |
closing_a_targets_flows_aborts_its_open_streams |
After close_flows, read, write and flush of a guarded stream fail |
supervisor/tests/unit/plane.rs |
closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive |
A receive already waiting fails, and so does the next send |
supervisor/tests/unit/plane.rs |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers |
Closing the balancer’s target leaves the stream; closing the member’s ends it |
supervisor/tests/unit/session.rs |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy |
The reason revocation precedes the table store |
supervisor/tests/unit/session.rs |
revoking_under_close_ends_only_the_sessions_of_that_principal |
Only the session bound to the revoked principal closes; the same user’s session on another inbound, which holds another principal, and another user’s session stay open |
supervisor/tests/unit/session.rs |
revoking_under_keep_ends_no_session |
The bound session stays open and may still admit flows; the principal is marked revoked under Keep |
supervisor/tests/unit/session.rs |
revoking_under_close_after_ends_the_session_once_the_grace_passes |
With a 10 s grace, the session is open 1 s before it and closed once it has passed |
supervisor/tests/unit/session.rs |
cancelling_the_root_cancels_every_session |
Bound or not, every session is a child of the root token |
supervisor/tests/unit/session.rs |
a_stopped_listener_opens_no_session |
The reason StopAccepting precedes CloseSessions |
app/tests/integration/e2e_reload.rs |
an_open_connection_keeps_transferring_across_a_reload |
Through the app: an open connection keeps its route across a reload that sends new ones to a blackhole |
app/tests/integration/e2e_reload.rs |
a_reload_that_removes_a_user_closes_only_their_connections |
Through the app: removing one user from the config closes only their connection |
app/tests/integration/e2e_hysteria_inbound.rs |
a_reload_keeps_the_inbound_serving_on_its_udp_port |
Through the app: a changed Hysteria 2 inbound on an unchanged bind keeps serving on its UDP port |
update and update_with, the DrainPolicy::CloseAfter timer firing, the edit_users errors (no such user set), and the spawn_tun fallback in Listener::swap have no dedicated test. Add one when you change those paths.