The supervisor
Source files: 54 · checked against Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/Cargo.tomlEtemenanki/supervisor/src/lib.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/policy.rsEtemenanki/supervisor/src/entity/mod.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/user.rsEtemenanki/supervisor/src/entity/usage.rsEtemenanki/supervisor/src/build/mod.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/build/users.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/build/outbound.rsEtemenanki/supervisor/src/topology/mod.rsEtemenanki/supervisor/src/topology/spec_plan/mod.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/router.rsEtemenanki/supervisor/src/topology/inbound/mod.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/system/mod.rsEtemenanki/supervisor/src/system/listener.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/sampler.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/socket_policy.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/supervisor/tests/support/mod.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/validate.rsEtemenanki/supervisor/benches/tracking.rsEtemenanki/environment/src/dial/socket.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/lower.rsEtemenanki/app/src/api.rsEtemenanki/ffi/src/proxy.rsEtemenanki/ffi/src/platform.rsEtemenanki/ffi/src/client.rsEtemenanki/webclient/src/tracking.rsEtemenanki/app/tests/unit/subscribe.rsEtemenanki/app/tests/integration/e2e_balancer.rsEtemenanki/app/tests/integration/e2e_dns.rskatana/src/manager/node.rskatana/src/lower/mod.rskatana/src/api/mod.rs
etemenanki-supervisor is the runtime every Etemenanki program runs on. etemenanki-app, the mobile library etemenanki-ffi and katana each lower their own input (a TOML file and its subscribe file, the strings a phone app passes in, a panel’s answers) into one typed Spec<U> and hand it to a Supervisor. The supervisor validates the spec, plans a reconcile from what runs, and applies it in two phases: a fallible prepare that changes nothing that runs, then an infallible commit. The reconcile keeps what did not change: a session does not belong to the listener that accepted it, an unchanged bind keeps its socket, and an unchanged outbound keeps its state. The crate knows nothing of config files, their formats, or how a reload is triggered.
This page is the entry to the supervisor section. It covers the split between the supervisor and its front ends, the module map, the public API, the actor every change goes through and the data flow of an apply, the cancellation tokens and the task tree, the socket policy, shutdown, and the ids that name what runs. The spec’s types are on The spec, the semantic rules on Validation, and every step of the plan, prepare and commit on Reconcile.
Responsibilities
Section titled “Responsibilities”The supervisor:
- validates a spec against every semantic rule and refuses it as a whole when one fails (
build/validate.rsis the only owner of those rules); - plans, per resource, whether to keep, build, swap, stop, close or drain, by comparing the running spec with the desired one;
- builds outbounds, balancers, the route table, the DNS resolvers and DNS service, inbound handlers and user tables;
- binds listeners: TCP, UDP for Hysteria 2, Unix sockets, and TUN devices it creates or is handed as a descriptor;
- runs the accept loops, registers every accepted connection as a session, and drives each connection’s protocol runtime;
- routes every new flow, mux sub-flow and UDP packet through the plane current at that moment;
- keeps a version per outbound build, and hands a replaced or removed version to its drain policy;
- admits users per inbound, revokes removed users under a removal policy, and enforces per-user speed limits;
- tracks every live flow, samples rates, and keeps a per-user usage ledger;
- stops all of it on shutdown.
It does not:
- parse any format or read config files. Certificates, keys and user lists arrive in the spec as bytes and typed values. The only files it reads are the geo data files a
RouteSpecnames, and the system hosts file when a split DNS spec setsuse_hosts. - decide when to reload. A file watcher, a panel poll or an API call is the front end’s business; the supervisor only answers
apply. - log config errors. It returns an
ApplyError, and the front end decides how to show it.
What stays live across an apply
Section titled “What stays live across an apply”The crate docs (supervisor/src/lib.rs) name four things an apply keeps:
| What | How it is kept | Owner page |
|---|---|---|
| Connections | Each accepted connection is a session under its own cancellation token, a child of the supervisor’s root token, not of its listener. | Users and sessions |
| Listeners | A listener is keyed by its BindSpec. An unchanged bind keeps its socket or device, and only the handler serving new connections is swapped. |
Serving |
| Routing | Routes and outbounds form one Plane behind an atomic cell. Every new flow, mux sub-flow and UDP packet reads the current plane; a route change does not re-route what is already open. |
Routing plane |
| Stateful outbounds | An outbound whose spec is unchanged, and whose resolvers did not change, is carried over by Arc, keeping its QUIC connection, tunnel or balancer health. A changed DnsSpec rebuilds every outbound except a blackhole (see Reconcile). |
Outbounds |
Users are a resource of their own. They change through set_users, upsert_user and remove_user, which rebuild only the user tables of the inbounds admitting the changed set. A user’s speed limit is shared by all their flows and takes effect on live flows at their next read or write, without ending their sessions.
These end live connections:
| Cause | Decided by | Default |
|---|---|---|
| A user no longer admitted by an inbound: removed from its set, or their credential for it changed | UserRemovalPolicy, per inbound or supervisor-wide |
Close |
| An outbound version replaced or removed | DrainPolicy, per outbound or supervisor-wide |
Keep: no action |
| An inbound removed from the spec, or its tag renamed | Always: its sessions are closed | none |
A change the planner marks Step::Disrupt: a TUN inbound’s device or settings, or a Hysteria 2 listener’s obfuscation |
Needs ApplyOptions::allow_disruptive; etemenanki-app refuses it with a restart hint |
refused |
close(selector) |
The caller | none |
Tracker::kill(id) or Tracker::kill_where(select): one flow, or the flows a predicate picks |
The caller: the FFI’s close_connection, the REST API (webclient/src/tracking.rs) |
none |
shutdown(grace) |
The caller; live connections get up to grace first |
none |
Policies
Section titled “Policies”supervisor/src/policy.rs holds the two policies and their supervisor-wide defaults:
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub enum UserRemovalPolicy { Keep, #[default] Close, CloseAfter(Duration),}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub enum DrainPolicy { #[default] Keep, Close, CloseAfter(Duration),}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct Policies { pub user_removal: UserRemovalPolicy, pub drain: DrainPolicy,}- A
UserRemovalPolicyis applied to each live session bound to a revoked principal (entity/session.rs→enforce):Keepleaves the session’s token alone,Closecancels it at once, andCloseAfter(grace)cancels it aftergraceunless the session ends first. - A
DrainPolicyis applied to an old outbound version (draininsupervisor.rs):Keeptakes no action on the old version’s flows,Closecancels the version’s closed token at once (Target::close_flows), andCloseAfter(grace)does so aftergrace. - The defaults differ on purpose. Removing a user defaults to
Closebecause revoking access is meant to take effect. Replacing an outbound defaults toKeepbecause a TCP flow cannot move to the successor mid-stream, so by default the drain does not close the old version’s flows. Spec::policiessets the supervisor-wide values.InboundSpec::user_removaloverrides the removal policy for one inbound (removal_policyinsupervisor.rsfalls back tospec.policies.user_removal), andOutboundSpec::drainoverrides the drain policy for one outbound (resolved by the planner, see Reconcile).- Both
CloseAftertimers run on the supervisor’sTaskTracker, so shutdown waits for them. A removal timer selects on its session’s token and a drain timer on the root token, and shutdown cancels the root (session tokens are its children), so both timers end at once instead of running out their grace (see The task tree).
What a front end owns
Section titled “What a front end owns”| Concern | Front end | Supervisor |
|---|---|---|
| Input format | Parses its own: TOML, a subscribe file, a panel’s JSON, FFI arguments | Takes only a Spec<U> |
| Parsing and lowering | Reports syntax errors and unknown keys; parses its strings (UUIDs, method names, addresses, keys, CIDRs); applies its schema’s defaults; names its users; refuses key combinations a spec cannot express | Never sees any of it |
| Semantics (duplicate tags, unknown references, how the typed pieces fit together) | Does not re-check them | build/validate.rs, the one owner |
| Files | Reads certificates, keys and user lists; passes bytes | Reads geo data and the hosts file only |
| User identity | Chooses the key type U |
Hands out an internal UserKey per admitted U |
| When to change | Watches files, polls a panel, serves an API | Answers apply, update and user edits |
| Disruptive changes | Chooses ApplyOptions |
Refuses them without allow_disruptive |
| Host socket policy | Passes SocketOptions to the builder |
Applies it to every outbound socket |
| Usage | Polls take_usage or installs a UsageSink |
Counts wire bytes per user |
| Logging | Sets up tracing and wraps ApplyError in its own text |
Logs through tracing under its own targets |
Log lines
Section titled “Log lines”The actor and the listeners it starts log these lines. What the accept loops and connections log is on Serving.
| Level | Target | Line | When |
|---|---|---|---|
| debug | etemenanki_supervisor::supervisor |
applied: {} built, {} reused, {} swapped, {} drained |
At the end of every commit, with the lengths of those four report fields |
| info | etemenanki_supervisor::system::listener |
inbound {tag} listening on {bind} |
Listener::start, for each listener a commit starts |
| info | etemenanki_supervisor::system::listener |
inbound {tag} owns tun device {name} |
A bind that created a TUN device, with the device’s name |
| info | etemenanki_supervisor::system::listener |
inbound {tag} serves a supplied tun device |
A bind that adopted a supplied descriptor |
| debug | etemenanki_supervisor::system::listener |
could not remove {path}: {e} |
A Unix listener drops and its socket file, still the one it created, cannot be removed (SocketFile::drop) |
The crate
Section titled “The crate”supervisor/Cargo.toml declares etemenanki-supervisor version 0.3.1, published only to a private Cargo registry. It depends on etemenanki-concepts 2.0.0, etemenanki-environment 3.0.0 and etemenanki-protocols 4.0.1 with the hysteria and tun features always on, so every supervisor can serve Hysteria 2 and TUN. Among the other dependencies, the manifest’s own comments say why three are there: scc is the live-flow registry, quinn only names the TLS material a Hysteria 2 listener is swapped to, and serde lets a UsageDelta reach a panel as it is (“No config types: this crate never parses a config format”).
The development dependencies are tokio with test-util (so grace periods can be tested against a paused clock), openssl for test certificates, proptest for the usage ledger’s property tests, criterion for the tracking bench, and socket2 for the socket type a socket-policy hook is shown.
Module map
Section titled “Module map”Directorysupervisor/
- Cargo.toml
Directorysrc/
- lib.rs crate docs, modules and re-exports
- supervisor.rs
Supervisor,SupervisorBuilder,check, and the actor - policy.rs
UserRemovalPolicy,DrainPolicy,Policies - serve.rs accept loops and serving one connection (crate-private)
- connector.rs
AppConnector, the connector every inbound dials through Directorybuild/ constructing what runs from a validated spec
- mod.rs
- apply.rs
ApplyReport,ApplyError - validate.rs every semantic rule of a spec
- dns.rs resolvers and the DNS service
- inbound.rs inbound handlers and their user tables
- outbound.rs one outbound handler per protocol
- route.rs the route table and its geo data
- users.rs admissions, principals and user keys
Directoryentity/
- mod.rs
- id.rs
OutboundId,Resource,SessionId,UserKey,FlowId,RuleId - user.rs
UserId,UserName,UserSpec,UserSet, credentials - session.rs the session registry
- usage.rs wire counters, the usage ledger,
UsageSink
Directorysystem/
- mod.rs
- listener.rs binding a
BindSpecand the task serving it
Directorytopology/
- mod.rs
Directoryspec_plan/
- mod.rs
Spec,ObfsSpec,Secret - inbound.rs inbound specs
- outbound.rs outbound specs
- transport.rs stream shapes and TLS specs
- route.rs
RouteSpec,BalancerSpec - dns.rs
DnsSpec - plan.rs
RunningState,Plan,Step,plan
- mod.rs
Directoryinbound/
- mod.rs
BindSpecand the inbound handlers
- mod.rs
Directoryoutbound/
- mod.rs
Outbound, one variant per kind - proxy.rs proxy protocol clients
- freedom.rs direct TCP and UDP
- udp_fanout.rs
FanOutLink, routing a UDP association per packet - dns.rs the DNS service as an outbound,
RoutedDialer
- mod.rs
- plane.rs
Plane,PlaneCell,Target - router.rs the compiled route table
- balancer.rs balancers and their health probes
- flow.rs
Flow,Principal,FlowContext
Directorytrack/
- mod.rs
Tracker,FlowEntry,FlowEvent - metered.rs
Metered,MeteredDatagram - pace.rs per-user token buckets
- sampler.rs the sampler and
StatsSnapshot
- mod.rs
Directorytests/
- hot_swap.rs applies against live connections
- socket_policy.rs the socket-policy hook
- tracking.rs flows, sessions and usage
Directorysupport/
- mod.rs shared specs, echo servers, protocol clients
Directoryunit/ unit tests, compiled into the crate
- …
Directorybenches/
- tracking.rs
lib.rs makes build, connector, entity, policy, system, topology and track public, keeps serve crate-private and supervisor private, and re-exports the entry points at the crate root:
pub use build::apply;pub use supervisor::{ApplyOptions, Selector, SessionInfo, Supervisor, SupervisorBuilder, check};Inside build, only apply and validate are public; the builders (dns, inbound, outbound, route, users) are crate-private. system has no public item: system::listener is crate-private. The unit tests live under tests/unit/ and are compiled into the modules they test through #[path] attributes (plan.rs, validate.rs, session.rs, plane.rs, router.rs, balancer.rs, track.rs, usage.rs, serve.rs and dns_outbound.rs).
| Module | Documented on |
|---|---|
topology::spec_plan (types), policy, entity::user (types) |
The spec |
build::validate, build::apply |
Validation |
topology::spec_plan::plan, the actor’s prepare and commit |
Reconcile |
serve, system::listener, topology::inbound, build::inbound |
Serving |
entity::user, entity::session, build::users |
Users and sessions |
connector, topology::plane, topology::router, topology::flow, build::route |
Routing plane |
topology::outbound, topology::balancer, build::outbound |
Outbounds |
build::dns, topology::outbound::dns, topology::spec_plan DNS types |
DNS |
track |
Tracking |
entity::usage |
Usage |
supervisor (the handle, the builder, the actor), entity::id |
this page |
Key types
Section titled “Key types”Supervisor<U>
Section titled “Supervisor<U>”#[derive(Clone)]pub struct Supervisor<U: UserId> { commands: mpsc::Sender<Command<U>>, plane: PlaneCell, tracker: Tracker, usage: Arc<UsageBook<U>>,}
impl<U: UserId> Supervisor<U> { pub fn builder() -> SupervisorBuilder<U>; pub async fn start(spec: Spec<U>) -> Result<(Self, ApplyReport), ApplyError>;
pub async fn apply(&self, spec: Spec<U>) -> Result<ApplyReport, ApplyError>; pub async fn apply_with(&self, spec: Spec<U>, options: ApplyOptions) -> Result<ApplyReport, ApplyError>; pub async fn update(&self, edit: impl FnOnce(&mut Spec<U>) + Send + 'static) -> Result<ApplyReport, ApplyError>; pub async fn update_with( &self, edit: impl FnOnce(&mut Spec<U>) + Send + 'static, options: ApplyOptions, ) -> Result<ApplyReport, ApplyError>;
pub async fn set_users(&self, set: &str, users: BTreeMap<U, UserSpec>) -> Result<(), ApplyError>; pub async fn upsert_user(&self, set: &str, id: U, user: UserSpec) -> Result<(), ApplyError>; pub async fn remove_user(&self, set: &str, id: U) -> Result<bool, ApplyError>;
pub async fn close(&self, selector: Selector<U>) -> usize; pub async fn sessions(&self) -> Vec<SessionInfo<U>>;
pub fn epoch(&self) -> u64; pub fn tracker(&self) -> &Tracker; pub fn take_usage(&self) -> Vec<UsageDelta<U>>; pub fn usage_snapshot(&self) -> Vec<UsageTotal<U>>;
pub async fn shutdown(&self, grace: Duration);}A Supervisor is a handle. Cloning it gives another handle to the same running supervisor: the command sender, the plane cell, the tracker and the usage book are all shared. When the last handle is dropped without shutdown, the actor sees its command channel close and stops everything at once (see Shutdown).
| Method | Through the actor | What it does |
|---|---|---|
builder |
no | A SupervisorBuilder with every default. |
start |
runs the first apply itself | Self::builder().start(spec). |
apply, apply_with |
Command::Apply |
Reconcile to spec. apply passes ApplyOptions::default(). |
update, update_with |
Command::Update |
Reconcile to the running spec as edit changes it. |
set_users |
Command::Users |
Replace every user of the set set. |
upsert_user |
Command::Users |
Add one user to set, or replace their UserSpec. |
remove_user |
Command::Users |
Remove one user from set; true if they were in it. |
close |
Command::Close |
Close the sessions selector names at once; returns how many this call closed. |
sessions |
Command::Sessions |
Every live session, by id. |
epoch |
no | How many planes have been published: self.plane.load().epoch(). |
tracker |
no | The Tracker: live flows and sessions, flow events, sampled rates. |
take_usage |
no | Every user’s wire bytes since the previous take. |
usage_snapshot |
no | Every user’s wire bytes since start, taken or not. |
shutdown |
Command::Shutdown |
Stop accepting, give live connections up to grace, close the rest, wait for them. |
apply replaces the whole spec, user sets included. update starts from the spec the actor holds, which includes every user edit committed since the last apply, so a front end that edits users one at a time and later changes one listener can use update without re-sending its users.
take_usage returns at most one UsageDelta per user and none for a user who moved nothing. It visits one account per active user. Its price is lag: a live session’s bytes are counted at each sampler tick, so what it moved since the last tick comes in a later take; a session that has ended is counted in full. Bytes a usage sink took are never returned by take_usage, and the other way round. The ledger itself is on Usage.
SupervisorBuilder<U>
Section titled “SupervisorBuilder<U>”pub struct SupervisorBuilder<U: UserId> { options: ApplyOptions, sample_interval: Duration, usage_sink: Option<PushUsage<U>>, socket: SocketOptions,}
impl<U: UserId> Default for SupervisorBuilder<U>;
impl<U: UserId> SupervisorBuilder<U> { pub fn options(self, options: ApplyOptions) -> Self; pub fn sample_interval(self, interval: Duration) -> Self; pub fn usage_sink(self, sink: impl UsageSink<U>, interval: Duration) -> Self; pub fn socket_options(self, socket: SocketOptions) -> Self; pub async fn start(self, spec: Spec<U>) -> Result<(Supervisor<U>, ApplyReport), ApplyError>;}The builder holds what outlives any one spec:
| Setting | Default | Effect |
|---|---|---|
options |
ApplyOptions::default() |
Options for the first apply. Nothing runs before it, so the first plan never contains a disruptive step, and the option has no effect on it. |
sample_interval |
DEFAULT_SAMPLE_INTERVAL, 1 s |
How often the sampler publishes a StatsSnapshot and reconciles the usage ledger. start panics with the sample interval must not be zero for Duration::ZERO. |
usage_sink |
none | Hand every user’s usage to sink once per interval, and a last batch at shutdown, after the last session has ended. |
socket_options |
SocketOptions::default() |
The policy every outbound socket is opened under, for the supervisor’s lifetime (see Socket policy). |
usage_sink stores a boxed closure (PushUsage<U>) that start turns into a task running push_usage. The sink runs on that task, so a slow sink never stalls a connection: it only delays its next batch, which then covers the longer interval. The task works like this:
- The first batch comes one
intervalafterstart(tokio::time::interval_at(now + interval, interval)), and a tick missed while the sink was slow is not made up in a burst (MissedTickBehavior::Delay). - On each tick, and once more when
background_stopis cancelled, it reconciles the ledger on the blocking pool (spawn_blocking(move || ledger.reconcile())), so the batch holds everything moved up to that moment rather than up to the last sampler tick. A panic in the reconcile is re-raised on the task withstd::panic::resume_unwind. - It then takes the usage (
UsageBook::take) and callssink.report(&batch)only when the batch is not empty, as theUsageSinkcontract promises (“Never called with an empty batch”).
How a batch is formed is on Usage.
The sampler (track/sampler.rs → Sampler::run) ticks first one sample_interval after start (interval_at), with MissedTickBehavior::Delay so a tick delayed by a busy runtime is not made up in a burst. Each tick runs on the blocking pool (spawn_blocking, a panic re-raised with resume_unwind), reconciles the usage ledger, and publishes its StatsSnapshot with watch::Sender::send_replace. When background_stop is cancelled it runs one final tick and returns. The tick itself is on Tracking.
start does, in order:
- Asserts that
sample_intervalis not zero. - Builds an
Actorwith the socket policy (Actor::new), which creates the root token, theTaskTracker, the usage ledger, the session registry, the tracker and its (not yet running) sampler, and the empty plane. - Runs the first apply directly on the caller’s task:
actor.apply(spec, options).await?. A refused spec returns itsApplyErrorhere. Nothing has been spawned yet, and any listener the refused prepare bound was released when prepare returned its error. - Spawns the sampler (
tokio::spawn(sampler.run(sample_interval, background_stop))). - Spawns the usage push task when a sink is set, under the same
background_stoptoken. - Creates the command channel,
mpsc::channel(16). - Builds the handle from the actor’s plane cell, tracker and usage book.
- Spawns the actor’s loop,
tokio::spawn(actor.run(rx)), and returns the handle with the firstApplyReport.
Listeners started by the first commit begin accepting in step 3, before the sampler runs. That is harmless: the sampler counts what flows moved whenever it first ticks.
ApplyOptions
Section titled “ApplyOptions”#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct ApplyOptions { pub allow_disruptive: bool,}The planner marks two kinds of change as Step::Disrupt, because each ends live connections: a change to a TUN inbound (its device or its settings), and a change to a Hysteria 2 listener’s obfuscation. Such a change needs allow_disruptive. Without the flag the actor refuses the whole spec after planning and before preparing anything:
inbound <tag>: <reason>; this ends its live connections and needs allow_disruptiveThe actor reports only the first disruption in plan order (plan.disruptions().next()), so ApplyError::Disruptive names one inbound even when several need the flag. The reasons and the steps are on Reconcile.
etemenanki-app never sets the flag, so it refuses such a reload with a restart hint, and restarting the app applies the change. katana always sets it, which lets a Hysteria 2 obfuscation change from the panel through. a_hysteria2_obfuscation_change_needs_allow_disruptive pins the Hysteria 2 case: refused without the flag; with it, restarted names the inbound and the port is kept.
Selector<U>
Section titled “Selector<U>”#[derive(Debug, Clone, PartialEq, Eq)]pub enum Selector<U> { Session(SessionId), User(U), Inbound(CompactString), All,}| Selector | Closes |
|---|---|
Session(id) |
That session, if it is live. |
User(id) |
Every session bound to that user, on any inbound. A user id no inbound has admitted has no key and closes nothing. A session whose first flow has not yet said who it is belongs to no user yet. |
Inbound(tag) |
Every session of that inbound, bound to a user or not. |
All |
Every live session. |
The actor maps the selector to a registry scope (Scope::One, Scope::User, Scope::Inbound, Scope::All), mapping a user id to its UserKey first. The count is of sessions this call cancelled: a session already closed is not counted again.
SessionInfo<U>
Section titled “SessionInfo<U>”#[derive(Debug, Clone, PartialEq, Eq)]pub struct SessionInfo<U> { pub id: SessionId, pub inbound: CompactString, pub source: Option<IpAddr>, pub user: Option<U>, pub up: u64, pub down: u64, pub started: Instant,}sessions builds these from the registry’s SessionStats, in session id order, translating each UserKey back to the front end’s U. user is None until the session’s first flow binds it to a user. It stays None for a session that never authenticates, and for a session admitted without a user, whose principal is Principal::anonymous() (for example on an inbound that names no user set). up and down are wire bytes (see Usage). started is a std::time::Instant.
Tracker::sessions returns the same registry as SessionStats (with a UserKey and the user’s label instead of U) without going through the actor. The REST API uses that path.
pub async fn check<U: UserId>(spec: &Spec<U>) -> Result<(), ApplyError>;check is a dry run of a first start. It builds a throwaway Actor with SocketOptions::default(), plans from the empty running state, and runs prepare with binding switched off (prepare(&plan, spec, false)), then drops everything it built.
check does |
check does not |
|---|---|
Run every validation rule (plan calls validate first) |
Bind a TCP, UDP or Unix listener |
| Build the DNS resolvers, reading the hosts file if the spec asks for it | Create or adopt a TUN device |
| Build every outbound, loading its TLS material | Pair a handler with a bound socket or device (Pending::pair), so no Hysteria 2 endpoint is opened and no TUN device is prepared for a runtime |
| Build every balancer | Spawn a task, start the sampler, or dial anything |
| Compile the route table, reading its geo data | Check for disruptive changes (nothing runs) |
| Admit users and build every inbound handler, parsing certificates and keys | Use the host’s socket policy (it never needs one) |
It refuses what a start would, short of binding and of what pairing a bound socket or device with its handler does: opening a Hysteria 2 endpoint (with its obfuscation) and preparing a TUN device for its runtime. etemenanki-app calls it through instance::check_bytes in two places: its --test flag (instance::check, over the file named by -c), and the REST API’s route switch (app/src/api.rs → set_route), which checks the edited config before writing it. Captured from the pinned binary (timestamps and colours removed):
$ etemenanki-app --test -c ok.tomlConfiguration OK.
$ etemenanki-app --test -c duplicate-tag.tomlERROR etemenanki_app: configuration invalid: duplicate outbound tag direct
$ etemenanki-app --test -c bad-ca.tomlERROR etemenanki_app: configuration invalid: building outbound proxy@v1 failed: error:04800064:PEM routines:PEM_read_bio_ex:bad base64 decode:../crypto/pem/pem_lib.c:968:The text after failed: is OpenSSL’s and varies with the OpenSSL version the binary links. The second error comes from validation. The third comes from building an outbound during prepare: the app read the CA file and passed its bytes in the spec, and the supervisor’s outbound builder, the same one a start uses, could not parse them. configuration invalid: is the app’s prefix; the rest is the ApplyError’s Display.
The user id type U
Section titled “The user id type U”Supervisor, its builder, and every public type that names a user are generic over the front end’s user key:
pub trait UserId: Clone + Eq + Ord + Hash + Debug + Send + Sync + 'static { fn name(&self) -> UserName;}U keys a UserSet and names the user in every UsageDelta, UsageTotal, SessionInfo and Selector::User. name gives the UserName (an email or a username, never a secret) that logs and the protocols’ user labels carry. The data plane never sees U: the actor hands out a UserKey per admitted id and connections carry that u64 (see Entity ids). etemenanki-app and the FFI key users by UserName itself; katana keys them by the panel’s numeric id. The trait and user sets are on Users and sessions.
ApplyReport and ApplyError
Section titled “ApplyReport and ApplyError”Both live in supervisor/src/build/apply.rs. An apply or update answers with an ApplyReport or an ApplyError; a user edit answers with () or a bool, or an ApplyError. ApplyReport lists what the apply did to each resource; every resource of the desired spec appears in reused or built. Commit fills it from the plan’s steps:
| Field | Filled from |
|---|---|
reused |
every Step::Reuse |
built |
every Step::Build |
swapped |
every Step::SwapHandler |
drained |
every Step::Drain. The drain policy runs only when the running target of that tag is the version the step names (target.id() == outbound), but the version is listed either way. |
rebound |
a Step::Bind for an inbound whose tag had a different bind in the running spec; a Step::Bind for a new tag is reported under built only |
restarted |
every Step::Disrupt |
removed |
every Step::CloseSessions |
The fields are explained on Reconcile.
ApplyError says why a change was not made, and a refused change leaves what runs as it was. Where each variant can arise:
| Variant | Display |
Raised in |
|---|---|---|
DuplicateTag |
duplicate {kind} tag {tag} |
plan → validate |
UnknownReference |
{from} references unknown {kind} {tag} |
plan → validate |
EmptyBalancer |
balancer {balancer} has no members |
plan → validate |
UnprobeableMember |
balancer {balancer}: outbound {member} has no upstream a TCP health probe can reach |
plan → validate |
Invalid |
{resource}: {reason} |
plan → validate; a user edit’s validate_admission; a user edit naming a missing set |
Disruptive |
inbound {inbound}: {reason}; this ends its live connections and needs allow_disruptive |
the actor, between plan and prepare |
Build |
building {resource} failed: {source} |
prepare; a user edit building an inbound’s user table or fitting it to the listener |
Bind |
inbound {inbound}: binding {bind} failed: {source} |
prepare, only when binding |
Stopped |
the supervisor has shut down |
the handle, when the actor is gone; the actor, for an update or a user edit with no running spec (unreachable after a successful start) |
Every rule and its exact reason text are on Validation.
The actor
Section titled “The actor”One task owns everything that runs. The handle talks to it over a channel, so applies and user changes are serialised without a lock held across an .await.
Commands
Section titled “Commands”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),}Every call that goes through the actor uses one private helper, ask:
- Create a
oneshotchannel for the reply. send(command).awaiton the boundedmpscchannel. With 16 commands already queued, the caller waits here for room.- Await the reply.
A failure at step 2 (the receiver is gone) or step 3 (the reply sender was dropped) becomes ApplyError::Stopped. close maps that to 0, sessions to an empty Vec, and shutdown ignores it.
The loop
Section titled “The loop”Actor::run receives one command at a time and finishes it before taking the next:
| Command | Handled by | Awaits inside |
|---|---|---|
Apply(spec, options, reply) |
apply |
creating a TUN device in prepare |
Update(edit, options, reply) |
clones the running spec, runs edit on it, then apply; with no running spec it answers Stopped |
creating a TUN device in prepare |
Users(set, edit, reply) |
edit_users |
nothing |
Close(selector, reply) |
close |
nothing |
Sessions(reply) |
sessions |
nothing |
Shutdown(grace, reply) |
shutdown, then replies and returns |
the grace, the tracked tasks, the background tasks |
When commands.recv() returns None, every handle has been dropped without a shutdown, and the loop runs shutdown(Duration::ZERO) before it returns.
Consequences of the single task:
- An apply and a user edit never interleave. No two applies prepare at once, so two prepares can never bind the same port against each other.
update’s edit runs on the actor task against the spec the actor holds, so nothing can land between reading the running spec and applying the edited one.closeandsessionswait behind a running apply. The only await inside an apply that can take time is in prepare, where a TUN device is created (BindSpec::Tun(TunSource::Create(..))→tun::open(spec).await); adopting a supplied descriptor, the other binds and every builder run synchronously on the actor task.- The actor holds no lock of its own. Its state is owned by the task and mutated through
&mut self. The locks it touches are held for synchronous sections only: the session registry’s mutex, the usage book’sRwLockand the tracker’s pacer map (allparking_lotlocks), a stream listener’swatchchannel (send_replace,borrow), and the Hysteria 2 inbound’s endpoint lock, whichset_quicandset_obfstake.
What bypasses the actor
Section titled “What bypasses the actor”Reads that must stay cheap never queue behind an apply:
| Read | Reads | Why it is safe without the actor |
|---|---|---|
epoch() |
the plane cell, ArcSwap::load |
The plane is published with one atomic store. |
tracker() |
the shared Tracker |
Flows register and deregister themselves; the registry is a concurrent map. |
take_usage(), usage_snapshot() |
the usage book | The ledger has its own locks: one mutex over its account maps (Ledger::inner, holding accounts and active) and one mutex per account. The key-to-id names are appended by the actor under a RwLock. |
All three keep working after shutdown.
The actor’s state
Section titled “The actor’s state”struct Actor<U: UserId> { state: RunningState<U>, shared: Shared, root: CancellationToken, dns: Option<Dns>, dns_target: Option<Arc<Target>>, targets: HashMap<CompactString, Arc<Target>>, probes: HashMap<CompactString, CancellationToken>, routes: Option<Arc<CompiledRoutes>>, listeners: Vec<Listener>, admissions: HashMap<CompactString, Admissions<U>>, keys: UserKeys<U>, usage: Arc<UsageBook<U>>, sampler: Option<Sampler>, background: Vec<JoinHandle<()>>, background_stop: CancellationToken, epoch: u64, socket: SocketOptions,}| Field | Holds |
|---|---|
state |
The spec last applied, and the latest version of every outbound and balancer tag, removed tags included (RunningState). |
shared |
What every accept loop shares (serve.rs → Shared): the plane cell, the session registry, the flow Tracker, and the TaskTracker. |
root |
The root cancellation token. |
dns, dns_target |
The resolvers built from the running DnsSpec, and the target the DNS service is reached through when the supervisor answers DNS itself. |
targets |
The current version of each outbound and balancer tag, as Arc<Target>. |
probes |
The probe token of each running balancer, by tag. |
routes |
The compiled route table. |
listeners |
Every bound listener with its serving task. |
admissions |
Who each inbound admits, by inbound tag: user id → principal and credential. |
keys |
The UserKey handed out for each admitted user id (build/users.rs → UserKeys). |
usage |
The usage book shared with every handle. |
sampler |
The sampler, until start spawns it. |
background, background_stop |
The sampler and usage push tasks, and the token that stops them. |
epoch |
Planes published so far. |
socket |
The socket policy every outbound socket is opened under. |
The usage book is the one piece of actor-built state a handle reads directly:
struct UsageBook<U> { ledger: Arc<Ledger>, names: RwLock<Vec<U>>,}names lists every id a key was handed out for, in key order; the actor appends the ids of new keys as it commits them (commit_keys). take and totals translate each ledger key back to U. A key with no name would mean a session bound a key before it was published, and taken bytes cannot be put back, so UsageBook::name panics with {key:?} was bound before it was published rather than drop the delta.
Data flow
Section titled “Data flow”An apply
Section titled “An apply”sequenceDiagram
participant FE as Front end
participant H as Supervisor handle
participant A as Actor task
participant R as What runs
FE->>H: apply_with(spec, options)
H->>A: Command::Apply on the channel
A->>A: plan(state, spec), validating first
alt the spec breaks a rule
A-->>H: Err(DuplicateTag, UnknownReference, EmptyBalancer, UnprobeableMember or Invalid)
else a Disrupt step without allow_disruptive
A-->>H: Err(Disruptive)
else valid, and no disruption or allowed
A->>A: prepare: build and bind, nothing that runs changes
alt prepare fails
A-->>H: Err(Build or Bind), what prepare built dropped
else prepare succeeds
A->>R: commit: keys, limits, plane, revocations, listeners, drains
A-->>H: Ok(ApplyReport)
end
end
H-->>FE: the result
-
Plan (
topology/spec_plan/plan.rs→plan) validates the spec and compares it withRunningState. It is pure: no sockets, no files, no clock. -
Disruption check (
Actor::apply): withoutallow_disruptive, the firstStep::Disruptrefuses the spec asApplyError::Disruptive. -
Prepare (
Actor::prepare) does everything that can fail, and changes nothing that runs: no listener starts, no plane is published, no user table or key changes. In order:- DNS: the running resolvers are kept unless the plan has
Step::Build(Resource::Dns). Rebuilding them (build_dns) also builds the DNS service’s target,Target::outbound(internal_id(self.epoch + 1), Outbound::Dns(service)), when the spec has the supervisor answer DNS itself; otherwise the runningdns_targetis kept. - Outbounds:
build_outboundfor eachStep::Build, the runningArc<Target>for eachStep::Reuse. - Balancers: for each
Step::Build, aMemberper member outbound with its probe target,Balancer::new, and aTarget::balancerunder the planned version; the new balancer is kept with itsprobe_intervalandprobe_timeoutfor commit. AStep::Reusekeeps the running target. - The route table: compiled again (
compile_routes) only onStep::Build(Resource::Route). - Admissions: for each inbound,
admiton a staged copy of the user keys, collecting the principals no longer admitted together with the inbound’s removal policy. - Inbounds, each in turn. An inbound the plan builds gets a new handler (
build_handler). A Hysteria 2 handler rebuilt on a running listener keeps that listener’s circuit semaphore whenmax_circuitsdid not change (same_circuits, which passesListener::circuits()to the builder). The handler is then prepared for the running listener on its bind (prepare_swap), or, with no listener on that bind, a new one is bound (listener::bind) and paired with it (Pending::pair);checkskips this last step. An inbound the plan keeps, whose listener is running, still gets a new user table (user_table, thenprepare_users) when its admissions differ from the running ones: another user, a replaced principal (Arc::ptr_eq) or another credential (same_admissions).
Prepare returns a
Preparedvalue (below). When a step fails, prepare returns its error, and what it built so far is dropped with it, listeners it bound included (a Unix listener’s socket file too). - DNS: the running resolvers are kept unless the plan has
-
Commit (
Actor::commit) cannot fail. It adopts the user keys (commit_keys), publishes speed limits, and replaces the targets. When the plan hasStep::PublishPlane, it publishes the plane (incrementingepoch), cancels the probe tokens of balancers the plan does not reuse, and starts the probes of new balancers under a child of the root, resolving with the new resolvers’ servers (dns.servers) and dialing withTcpDialer::new(self.socket.clone()); without a publish, probe tokens are left alone. It stores the route table, the resolvers and the DNS service target either way. It then revokes removed principals, starts new listeners (each withroot.child_token()as its stop token), swaps handlers, stores user tables, and walks the steps:Step::StopAcceptingstops the listener on that bind,Step::CloseSessionscloses the inbound’s sessions, andStep::Drainhands the old version to its drain policy. It ends by storing the admissions and the running state, and logsapplied: {} built, {} reused, {} swapped, {} drainedat debug level.
Every step, in order, is on Reconcile.
What prepare builds
Section titled “What prepare builds”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)>,}| Field | Holds |
|---|---|
dns, dns_target |
The resolvers and the DNS service target, rebuilt or carried over |
targets |
Every outbound and balancer tag of the spec, by tag |
new_balancers |
Each built balancer with its tag, probe_interval and probe_timeout, whose probes commit starts |
routes |
The route table, compiled or carried over |
admissions, keys |
Who each inbound admits, and the staged user keys |
revoked |
Per inbound, the principals no longer admitted and the removal policy that applies to them |
started |
New binds, each paired with its handler (Pending) |
swaps |
Handlers prepared for running listeners, by listener index |
users |
User tables prepared for running listeners, by listener index |
A Prepared that is dropped without a commit, as check does with map(drop), releases what it holds the same way.
Speed limits
Section titled “Speed limits”Actor::publish_speed_limits runs at the start of every commit and every user edit’s commit, never in prepare, so a refused change leaves live limits alone. It walks every user of every set in the spec and collects each user’s speed_limit by UserKey. A user in several sets gets the smallest of their limits, and a user with no key (one no inbound has admitted) is skipped. Tracker::set_speed_limits then gives each listed user’s pacer its rate and sets every other pacer it holds back to unlimited, so live flows feel the change on their next poll. The pacers are on Tracking.
A user edit
Section titled “A user edit”set_users, upsert_user and remove_user do not plan. Actor::edit_users changes one user set in a copy of the running spec and rebuilds only the user tables of the inbounds admitting that set:
-
Take the running spec (with none, the edit fails as
ApplyError::Stopped) and find the set. A missing set isApplyError::InvalidforResource::UserSet, shown asuser set {tag}: no such user set. -
Apply the edit. Removing a user who is not in the set returns
Ok(false)and changes nothing.SetandUpsertalways count as a change. -
Prepare, for each inbound whose
usersnames the set:- check the set against the inbound with
validate_admission, the only validation rule a user edit runs; - admit its users on a staged copy of the keys;
- find the inbound’s listener by bind, with
expect("every running inbound has its listener"); - build the user table (
user_table) and check that it is for the protocol the listener’s current handler serves (Listener::prepare_users→UserTable::fits). A mismatch fails the edit asApplyError::Buildwiththe handler is not of the kind its listener serves(io::ErrorKind::InvalidInput).
The first failure refuses the whole edit, and no table is stored.
- check the set against the inbound with
-
Commit: adopt the keys, publish speed limits, revoke the principals that are gone under each inbound’s removal policy, then store each table and its admissions, and keep the edited spec as the running one.
A user edit never publishes a plane, so epoch does not move. As in an apply, the principals of removed users are revoked before any new user table is stored, and a session whose first flow presents a revoked principal is closed under every policy (Session::admit). The details are on Users and sessions.
Tokens, tasks and the task tree
Section titled “Tokens, tasks and the task tree”Cancellation tokens
Section titled “Cancellation tokens”flowchart TB root["root token (Actor::new)"] stop["listener stop token, one per listener"] tun["TUN runtime token, one per runtime"] probe["probe token, one per new balancer"] session["session token, one per session"] bg["background_stop (independent)"] closed["Target closed token, one per outbound version (independent)"] root --> stop stop --> tun root --> probe root --> session
| Token | Created by | Parent | Cancelled by | Stops |
|---|---|---|---|---|
root |
Actor::new |
none | shutdown, once the grace is over |
Everything below it; drain grace timers; a Hysteria 2 endpoint still draining |
Listener stop |
commit: root.child_token() per started listener |
root |
Step::StopAccepting, shutdown, or the root |
The accept loop; Sessions::open refuses new sessions for it; the TUN runtime under it |
| TUN runtime token | run_tun_inbound: stop.child_token() per runtime |
its listener’s stop |
A restart sent by Listener::swap, its listener’s stop, or its listener handle being dropped (the restart channel closes) |
One TUN runtime and every flow it serves |
| Probe token | commit: root.child_token() per newly built balancer |
root |
A publish that drops or rebuilds the balancer, shutdown, or the root |
That balancer’s probe tasks |
| Session token | Sessions::open: root.child_token() |
root |
close, a removal policy, Step::CloseSessions, Session::admit refusing a revoked principal, the session’s own end (Session::drop), or the root |
On a stream listener, the tasks serving that connection; a Hysteria 2 connection, which the protocol closes; a removal grace timer waiting on it |
background_stop |
Actor::new: CancellationToken::new() |
none | shutdown, after every tracked task has ended |
The sampler and the usage push task |
Target closed token |
Target::outbound and Target::balancer: CancellationToken::new() |
none | Target::close_flows, under a drain policy of Close or CloseAfter |
Every stream and datagram link opened on that outbound version (Guarded). A flow sent through a balancer is guarded by the token of the member that carries it. |
Session tokens hang off the root, not off the listener that accepted the connection. That is what lets an apply stop, rebind or re-handle a listener without touching its connections: stopping a listener stops accepting and nothing else. What does end sessions is listed under What stays live across an apply.
A session ends when the last Arc<Session> drops. Session::drop cancels its token, which ends a removal grace timer still waiting on it; closes the session in its user’s ledger account, folding in the bytes not yet counted (Account::close); and removes it from the registry. Session ids count up from 1 (Sessions::next, AtomicU64::new(1)).
background_stop is deliberately not a child of the root. The sampler and the usage sink must outlive every connection so they count its last bytes, and cancelling the root is what ends the connections.
The task tree
Section titled “The task tree”The supervisor runs its tasks in three places: plain tokio::spawn, its one TaskTracker, and the blocking pool.
| Task | Spawned with | Spawned by | Ends when |
|---|---|---|---|
| The actor | tokio::spawn |
SupervisorBuilder::start |
Shutdown has been handled, or every handle is gone |
| The sampler | tokio::spawn, kept in background |
SupervisorBuilder::start |
background_stop, after one last tick |
| The usage push | tokio::spawn, kept in background |
SupervisorBuilder::start, with a sink |
background_stop, after one last reconcile and take, reported if not empty |
| A stream accept loop | TaskTracker::spawn |
Listener::start |
its listener’s stop. The select! is biased toward the stop, and the loop also leaves when Sessions::open finds the stop cancelled. |
| A Hysteria 2 listener | TaskTracker::spawn |
Listener::start |
its stop and its last connection, or the root |
| A TUN runtime loop | TaskTracker::spawn |
Listener::start, and a swap whose runtime already ended |
its listener’s stop, the listener handle being dropped, or its runtime ending on its own |
| One accepted socket | TaskTracker::spawn under select! with the socket’s session token (Shared::spawn_until) |
the stream accept loop | its transport ends, or its session token is cancelled. It runs the transport; a Unix socket’s connection is served on this task. |
| One byte stream a TCP transport yields | Shared::spawn_until under its session’s token |
the socket’s task, from the transport’s accept callback | the stream’s runtime ends, or its session token is cancelled. The session is the socket’s own, or, on a gRPC carrier, one opened for the stream. |
| One Hysteria 2 connection | a JoinSet inside its listener’s task (Hy2Inbound::run) |
the Hysteria 2 listener | the connection ends, or the protocol closes it when its session token is cancelled |
| A balancer probe, one per member | TaskTracker::spawn |
Balancer::spawn_probe, from commit |
its probe token |
| A removal grace timer | TaskTracker::spawn |
Sessions::revoke under CloseAfter |
the grace passes, or the session token is cancelled |
| A drain grace timer | TaskTracker::spawn |
drain under CloseAfter |
the grace passes, or the root is cancelled |
| A sampler tick, a ledger reconcile | spawn_blocking |
the sampler, the usage push | the tick or reconcile is done |
There is one TaskTracker, created in Actor::new and shared by Shared and the session registry, so shutdown can wait for every accept loop, connection, probe and timer with one wait; a Hysteria 2 listener’s task waits for its own connections. Shared::spawn_until wraps each socket and each stream as
self.tracker.spawn(async move { tokio::select! { _ = token.cancelled() => {} _ = fut => {} }});so cancelling a session token drops the future, and its sockets with it, at its next poll. On a gRPC carrier each HTTP/2 stream is a session of its own. A Hysteria 2 connection that finishes its QUIC handshake after its listener’s stop was cancelled gets no session: run_hysteria_inbound hands it a token that is already cancelled, so no session escapes the close of a removed inbound. The accept loops and what they serve are on Serving.
Socket policy
Section titled “Socket policy”pub fn socket_options(self, socket: SocketOptions) -> Self;SocketOptions (environment/src/dial/socket.rs) is the policy applied to every socket a dialer opens: a source bind_address (applied only to sockets of its own family), an interface, a packet mark (SO_MARK), a tcp_keepalive idle time, and a hook run on each socket after it is created and before it is bound or connected. The fields and platform support are on Dialers.
The builder stores the policy in the actor, which passes it to every builder that opens an outbound socket:
| Socket | Receives the policy through |
|---|---|
| Every dial of a proxy outbound, over TCP and its transports | build/outbound.rs → Dialer::new(socket) |
| A direct (freedom) flow’s TCP dial and UDP sockets | build/outbound.rs → Dialer::new(socket), whose TCP and UDP halves both carry it |
| A SOCKS outbound’s UDP relay socket | build/outbound.rs → UdpDialer::new(socket) |
| A Hysteria 2 outbound’s QUIC socket | Hy2Config::socket |
| A WireGuard outbound’s tunnel socket | WgConfig::socket |
| Every query a resolver sends from this host | build/dns.rs → ResolverOptions::socket |
| Balancer health probes | commit → TcpDialer::new(self.socket.clone()) |
The policy holds for the supervisor’s lifetime; an apply cannot change it. It does not reach:
- listener sockets, which
system::listener::bindopens with the standard library: the policy is for traffic this process originates; - the system resolver.
getaddrinfoopens its sockets inside the C library, where no policy reaches them. A resolver whose backend is the system resolver (for exampleDnsSpec::Single(Backend::System)) uses it.
This is where a mobile VPN keeps the proxy’s own traffic out of its tunnel. etemenanki-ffi passes a hook that calls the app’s Platform::protect(fd) (VpnService.protect on Android); a socket the platform refuses fails with the platform refused to protect the socket, and the dial it was for fails with it. A hook error of any kind fails the dial: nothing leaves on a socket the policy did not accept.
check builds with SocketOptions::default(): it never dials, so it needs no policy.
Shutdown
Section titled “Shutdown”pub async fn shutdown(&self, grace: Duration);shutdown sends Command::Shutdown(grace, reply), and the actor runs Actor::shutdown:
- Cancel every listener’s
stoptoken. Accept loops return and release their sockets (a Unix listener removes its socket file);Sessions::openrefuses any session for them. A Hysteria 2 listener stops admitting and drains. A TUN runtime’s token is a child of its listener’sstop, so every TUN runtime and the flows it serves end here, before the grace: TUN flows are not sessions. - Cancel every balancer’s probe token.
- Close the
TaskTracker, sowaitcan finish once it is empty. - Wait up to
gracefor every tracked task to end on its own:tokio::time::timeout(grace, tracker.wait()). - Cancel the root. Every session token is a child, so every task still serving a stream connection ends at its next poll and the protocol closes every Hysteria 2 connection still open; removal and drain grace timers end; a Hysteria 2 listener closes its endpoint under the draining run, so a client still in its QUIC handshake (which has no session yet) ends too.
- Wait for every tracked task to end.
- Drop every
Listenerhandle. - Cancel
background_stopand await the sampler and the usage push task. Every connection has ended and folded its last bytes in, so the sampler runs its last tick, and the push task reconciles and takes once more and reports that batch if it is not empty.
The actor then replies and returns, dropping the command receiver. grace is how long live connections may take to finish; Duration::ZERO closes them at once.
When every handle is dropped without a shutdown, the loop runs the same sequence with Duration::ZERO. A Tracker clone taken with tracker() does not keep the actor alive.
After shutdown:
| Call | Result |
|---|---|
apply, apply_with, update, update_with |
Err(ApplyError::Stopped): the supervisor has shut down |
set_users, upsert_user, remove_user |
Err(ApplyError::Stopped) |
close |
0 |
sessions |
an empty Vec |
shutdown |
returns at once |
epoch |
the last published epoch |
tracker |
still readable; the last StatsSnapshot is the final tick’s |
take_usage, usage_snapshot |
still work: what no sink took can still be taken |
A command queued behind Shutdown gets Stopped too: its reply sender is dropped with the channel.
Entity ids
Section titled “Entity ids”supervisor/src/entity/id.rs names what runs by value instead of by pointer. A fan-out UDP link keys each sub-link by the OutboundId it was opened on rather than by the address of an Arc, because 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.
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]pub struct OutboundId { pub tag: CompactString, pub version: u64,}
pub enum ResourceKind { Inbound, Outbound, Balancer, UserSet }
pub enum Resource { Inbound(CompactString), Outbound(OutboundId), Balancer(CompactString), UserSet(CompactString), Route, Dns,}
pub struct SessionId(u64);pub struct UserKey(u64);pub struct FlowId(u64);pub struct RuleId(u32);| Id | Names | Assigned by | First value | Display |
|---|---|---|---|---|
OutboundId |
One build of one outbound or balancer tag | plan: 1 for a tag it has no version for, the previous version plus one on each rebuild |
1 |
{tag}@v{version} |
ResourceKind |
A kind of tagged resource, in DuplicateTag and UnknownReference |
— | — | inbound, outbound, balancer, user set |
Resource |
One resource, in a Plan, an ApplyReport and an ApplyError |
— | — | inbound {tag}, outbound {tag}@v{n}, balancer {tag}, user set {tag}, route, dns |
SessionId |
One accepted connection: a socket, a QUIC connection, or one stream of a gRPC carrier | Sessions::open, an AtomicU64 counter |
1 |
session#{n} |
UserKey |
One user id an inbound has admitted | UserKeys::key, called by admit the first time an inbound admits that id |
1 |
none |
FlowId |
One routed flow: a stream, or one sub-link of a UDP association | the Tracker, an AtomicU64 counter |
1 |
flow#{n} |
RuleId |
The index of a rule in RouteSpec::rules of the plane that routed a flow |
CompiledRoutes, per rule |
0 |
none |
- Versions never go back.
RunningState::versionskeeps the version of a removed tag, so a tag added back gets a version its old build, which may still be draining, never had. Balancers are versioned in the same map, and a balancer target carries anOutboundIdtoo, althoughResource::Balancerholds only the tag. - The empty tag is the supervisor’s own. Validation refuses an empty tag of any kind (
a {kind} tag must not be empty), so the supervisor uses it for the targets it makes itself (internal_id(version)insupervisor.rsbuilds anOutboundIdwith the empty tag): the empty plane’s blackhole (version0, before the first apply) and the DNS service target. That target is built in prepare only when the DNS is rebuilt, with the epoch of the plane about to carry it (internal_id(epoch + 1)); otherwise the running one is kept. The sampler reports the DNS service under the empty outbound tag. - Session and flow ids are unique for the supervisor’s lifetime, never reused.
SessionId,UserKey,FlowIdandRuleIdexposenew(raw)andget()asconst fn, so a front end can hand an id out and take it back, as the FFI does withFlowId::new(id)forclose_connection. - A user’s key stays assigned after the user is removed, so sessions a
Keeppolicy spared can still be found bySelector::User, and their last usage is still reported under the user’s id. A key is handed out only when an inbound admits the user, that is, when the user has the credential kind the inbound’s protocol takes; a user no inbound admits has no key. Keynnames thenth id handed out (id_of:ids[n - 1]);UserKey(0)names nobody. - A rule id belongs to one plane. A flow’s rule (
Routed::rule,FlowEntry::rule) isNonefor the default route and for a flow intercepted for the DNS service. The router turns the index into au32withexpect("fewer than 2^32 rules").
How the front ends use it
Section titled “How the front ends use it”| etemenanki-app | etemenanki-ffi | katana | |
|---|---|---|---|
| Code | app/src/instance.rs → Core, Instance |
ffi/src/proxy.rs → Proxy, running the app’s Core |
src/manager/node.rs → NodeManager |
U |
UserName (app/src/lower.rs → UserKey) |
UserName, through the app’s lowering |
Uid(i64), the panel’s user id; name() is UserName::Username of the number |
| Supervisors | One per process | One per started proxy | One per node |
| Start | Supervisor::builder() with every default |
Supervisor::builder().socket_options(...) with the protect hook |
Supervisor::start(spec) with every default |
| Changes | apply(spec) on each reload |
apply(spec) through the app’s Core::reload_with |
apply_with(spec, ApplyOptions { allow_disruptive: true }), with a spec lowered from the panel’s answers |
| Reads | tracker() for the REST API |
tracker(): stats(), flows(), kill(FlowId) |
take_usage() for traffic reports, tracker().subscribe() for audit |
| Shutdown | shutdown(SHUTDOWN_GRACE), 5 s, on a signal |
shutdown(STOP_GRACE), 2 s, in Proxy::stop; then RUNTIME_GRACE, 3 s, for the runtime’s threads |
shutdown(Duration::ZERO) when the node stops, then one last take_usage for its final report |
- etemenanki-app lowers the TOML config and its subscribe file into a spec (
lower), runscheckbehind--testand before the REST API’s route switch writes the edited config, and applies each reload with default options. It logsconfig loaded: {summary}after the first apply andconfig reloaded: {summary}after each applied reload, where the summary lists each non-empty group of theApplyReport(built [...]; swapped [...], and so on) or readsnothing to run. A disruptive change is refused and logged asreload refused, keeping the running config: inbound {inbound}: {reason}, which would end its live connections; restart to apply it. Its lowering is described on etemenanki-app: from TOML to a spec, and its reload on etemenanki-app: running, reloading and shutting down. When the configured REST API cannot bind at start, the binary shuts the supervisor down withDuration::ZEROand exits. - etemenanki-ffi runs the app’s
Corewith a lowering of its own (ffi/src/client.rs→lowering). It refuses an inbound that serves a proxy protocol to others: every server protocol, and a SOCKS or HTTP inbound that is not on a loopback address or a Unix socket (ServerInbound). It refuses more than one TUN inbound, binds the TUN inbound to the descriptor the phone supplies as aBindSpec::Tun(TunSource::Fd(...)), and adds one taggedtun(TUN_TAG) with the protocol’s defaults (DEFAULT_MTU,DEFAULT_UDP_IDLE_TIMEOUT,DEFAULT_MAX_FLOWS, sniffing on) when the config has none. It does not serve the config’s[api]. Its socket policy is the one on this page. It is described on Mobile library (etemenanki-ffi). - katana runs one supervisor per node and applies a spec lowered from the panel’s node and user list , with
allow_disruptiveset. It then callsclose(Selector::All)when the node’s transport or protocol settings differ from the ones last applied (NodeInfo::transport_eq,NodeInfo::protocol_eq, which cover among others the port, TLS, obfuscation, node type, cipher and server key), or when a local config edit the panel’s answer cannot show asked for it. Its audit task subscribes to flow events withtracker().subscribe()and records audit hits fromFlowEvent::Opened. See Node manager and Traffic accounting.
Design rules
Section titled “Design rules”The crate follows a few rules that the rest of this section relies on. The general principles are on Design principles.
- One owner per rule. Every rule about how the typed pieces of a spec fit together lives in
build/validate.rs. A front end owns what only it can know: parsing its strings, reading its files, the defaults of its schema, how it names users, and key combinations a spec could not express. Asapp/src/lower.rsputs it, a rule checked in a front end as well would be a second copy that drifts, and a rule checked only there would be missing for every other front end. - A typed spec compared with
==. ASpecholds typed values (aDestination, aUuid, decoded key bytes, PEM bytes) rather than raw config text. That lets the planner decide what to keep and what to rebuild by equality, without knowing where either spec came from. - Prepare, then commit. Everything fallible happens before anything that runs changes, so a refused spec leaves the running state exactly as it was.
- Read-copy-update routing. Routes and outbounds are published together with one atomic store, so a rule can never name a tag the plane lacks, and an old plane is freed when the last flow holding one of its pieces ends.
- Identities instead of pointers. Outbound versions, sessions, users, flows and rules are named by value (see Entity ids).
- A slow observer does not hold up the data plane. Each flow’s counters are its own; the sampler aggregates them once per tick on the blocking pool; flow events go out on a bounded broadcast channel (
EVENT_CAPACITY, 1024), where a subscriber that falls behind loses events rather than holding up the sender; a usage sink runs on its own task. Nothing per byte goes through the actor.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| A refused spec changes nothing: listeners, the plane and users stay as they were | Prepare changes nothing that runs; commit cannot fail; listeners bound by a refused prepare are released when it returns its error | a_refused_spec_changes_nothing (a validation refusal and a failed bind) |
| Applies and user edits never run concurrently | One actor task, one command at a time | none directly |
| A connection survives an apply that keeps its inbound | Session tokens are children of the root; an unchanged bind keeps its listener | an_established_connection_survives_an_apply_that_keeps_its_inbound |
| A route change reaches the next flow, mux sub-flow and UDP packet, and does not re-route an open TCP flow | The connector loads the plane once per connect; the fan-out routes per packet |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow, a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
epoch counts published planes; a refused apply publishes none |
epoch += 1 only at Step::PublishPlane, in commit |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow (epoch + 1), a_refused_spec_changes_nothing (unchanged) |
| The plane before the first apply drops everything | Plane::empty, epoch 0 |
the_empty_plane_drops_everything |
| An outbound whose spec and resolvers are unchanged is carried over; a rebuilt one is a new version | Step::Reuse inserts the running Arc<Target>; Step::Build takes the next version; a changed DnsSpec rebuilds every non-blackhole outbound |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply, each_rebuild_of_a_tag_takes_the_next_version, a_tag_added_back_takes_a_version_it_never_had |
A Step::Disrupt change needs allow_disruptive |
Plan::disruptions checked before prepare |
a_hysteria2_obfuscation_change_needs_allow_disruptive, a_tun_device_change_needs_allow_disruptive |
| A protocol change on the same bind keeps the socket | Listener keyed by BindSpec; handler swap |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
| Every session is a child of the root | Sessions::open → root.child_token() |
cancelling_the_root_cancels_every_session |
| A stopped listener opens no session | Sessions::open checks stop under the registry lock |
a_stopped_listener_opens_no_session |
| Removing an inbound closes its sessions | Step::StopAccepting, then Step::CloseSessions |
removing_an_inbound_closes_its_sessions |
close counts only what it closed |
Sessions::close skips cancelled tokens |
a_session_already_closed_is_not_counted_again |
| Shutdown does not wait out a grace timer | Timers run on the TaskTracker; removal timers select on the session token and drain timers on the root, and shutdown cancels the root |
shutdown_ends_pending_grace_timers |
| Shutdown does not wait for a stalled Hysteria 2 handshake | run_hysteria_inbound closes the endpoint once the root is cancelled |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
| The last usage batch is taken after the last session ended | background_stop is cancelled after tracker.wait() |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
| A slow sink never stalls a connection | The sink runs on its own task | a_slow_usage_sink_does_not_stall_the_data_plane |
| Every outbound socket is opened under the socket policy | The actor passes SocketOptions to every builder and to the probes |
direct_flows_are_dialed_on_hooked_sockets, a_socks_upstream_is_dialed_on_hooked_sockets, a_refusing_hook_fails_the_dial |
| A user key is published before any session can bind it | commit_keys runs first in commit and in a user edit’s commit |
none; UsageBook::name panics otherwise |
| Speed limits change only when an apply or user edit commits | publish_speed_limits is called only in commit |
a_speed_limit_change_keeps_the_users_sessions (the session is kept) |
| No spec tag is empty, so internal targets never collide with a spec’s | unique_tags in validate |
no_tag_may_be_empty |
Failure paths and cancellation
Section titled “Failure paths and cancellation”Starting
Section titled “Starting”start fails only through its first apply, with the errors a later apply could return, except Disruptive (nothing runs yet) and Stopped. Nothing has been spawned by then: the sampler, the usage push and the actor are spawned only after the first apply commits. Listeners bound during the refused prepare are released when prepare returns its error. start asserts that sample_interval is not zero (the sample interval must not be zero).
Refused changes
Section titled “Refused changes”A refused apply or user edit returns its ApplyError and leaves the running spec, versions, plane, listeners, user tables, keys and speed limits as they were. User keys handed out while preparing live on a staged copy (self.keys.clone()) and are adopted only at commit, so a refused change hands out no key.
A caller that stops waiting
Section titled “A caller that stops waiting”Dropping the future of a call does not cancel the command once it has been sent: the actor finishes it and its reply is discarded (let _ = reply.send(...)). An apply whose caller gave up still commits. A future dropped while it waits for room in the channel sends nothing.
Connections ended from outside
Section titled “Connections ended from outside”close, a removal policy, Step::CloseSessions and shutdown all end a connection through its session token. On a stream listener every task serving the connection runs under that token through Shared::spawn_until, so cancelling it drops the task’s future at its next poll. A Hysteria 2 connection is closed by the protocol when its session token is cancelled (see Serving).
Tracker::kill and Tracker::kill_where work per flow: they set the flow’s kill flag and wake it, so its next poll fails in either direction; the session it belongs to is untouched. A drain policy works one level down from sessions: it cancels the outbound version’s closed token, and every Guarded stream or datagram link opened on that version fails its next read or write with the outbound this flow was opened on was drained (ConnectionAborted), so the runtime relaying it ends that flow. The routing side is on Routing plane.
After shutdown
Section titled “After shutdown”Every call through the actor returns Stopped, 0, an empty list or nothing, as listed under Shutdown. The reads that bypass the actor keep answering.
Limits
Section titled “Limits”| Limit | Value | Where | Meaning |
|---|---|---|---|
| Command channel | 16 commands | SupervisorBuilder::start → mpsc::channel(16) |
A caller’s send waits while 16 commands are queued |
DEFAULT_SAMPLE_INTERVAL |
1 s | track/sampler.rs |
The sampler’s tick, and so the lag of take_usage for live sessions |
EVENT_CAPACITY |
1024 events | track/mod.rs |
Flow events a subscriber may fall behind before it loses some |
sample_interval |
must not be zero | assert in SupervisorBuilder::start |
Panics otherwise |
| Outbound versions | u64, from 1 per tag |
plan |
Kept for removed tags, so RunningState::versions has one entry per tag ever seen |
SessionId, FlowId |
u64, from 1 |
Sessions::open, Tracker |
Unique for the supervisor’s lifetime |
UserKey |
u64, from 1 |
UserKeys::key |
One per user id, assigned the first time an inbound admits it |
RuleId |
u32 |
topology/router.rs |
Fewer than 2^32 rules per route table |
| etemenanki-app shutdown grace | 5 s | app/src/main.rs → SHUTDOWN_GRACE |
|
| etemenanki-ffi stop grace | 2 s | ffi/src/proxy.rs → STOP_GRACE |
|
| etemenanki-ffi runtime grace | 3 s | ffi/src/proxy.rs → RUNTIME_GRACE |
After the stop grace, for the runtime’s threads (shutdown_timeout) |
| katana shutdown grace | 0 | src/manager/node.rs |
The per-inbound connection limits and the accept-error backoff belong to the serving layer; see Serving. The limits of the whole system are collected on Limits, timeouts and memory.
The integration tests start real supervisors on loopback sockets and drive them with minimal protocol clients from supervisor/tests/support/mod.rs:
| Helper | What it does |
|---|---|
start_on, start_built |
Start a supervisor on a spec built for freshly picked ports; on ApplyError::Bind pick again, up to 5 retries, since a port can be taken between picking and binding. start_built takes a builder factory, which is how the socket-policy and usage-sink tests install theirs. |
base_spec |
A direct freedom outbound, a block blackhole, the route defaulting to direct, DnsSpec::Single(Backend::System) and default policies. |
QUIET, SOON |
400 ms to wait for something that should not happen, 5 s for something that should. |
| Clients | SOCKS connect and UDP associate, HTTP CONNECT, a VLESS mux client, echo servers, a self-signed certificate. |
The hot-swap and tracking tests end with shutdown(Duration::ZERO) on every supervisor they started. The socket-policy tests drop their handles instead, which runs the same shutdown with a zero grace.
| Test | File | Behaviour it pins |
|---|---|---|
an_established_connection_survives_an_apply_that_keeps_its_inbound |
tests/hot_swap.rs |
A kept inbound is reused, not swapped or rebound, and its connection keeps transferring |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
tests/hot_swap.rs |
epoch moves by one; the next UDP packet and a new TCP flow follow the new route; the open TCP flow does not |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
tests/hot_swap.rs |
A mux sub-flow opened after the change takes the new route; the earlier one keeps its own |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close |
tests/hot_swap.rs |
On the freedom outbound direct: drained names the old version; under Keep the flow opened on it keeps echoing, under Close it is closed; new flows use the new version |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
tests/hot_swap.rs |
A reused Hysteria 2 outbound does not handshake again; a rebuilt one dials anew while the flow on the old one keeps echoing |
a_hysteria2_obfuscation_change_needs_allow_disruptive |
tests/hot_swap.rs |
Disruptive without the flag; with it, restarted names the inbound and the port is kept |
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections |
tests/hot_swap.rs |
swapped, not rebound; a connection accepted before the swap is served by the old handler |
a_refused_spec_changes_nothing |
tests/hot_swap.rs |
An unknown reference and a failed bind leave the epoch, the users and the routes; a listener bound during the refused prepare is released |
removing_a_user_closes_only_their_sessions_and_refuses_them_after |
tests/hot_swap.rs |
On a SOCKS inbound, under Close and Keep: the removed user’s session is closed or kept by the policy, the other user’s is untouched, and a new login by the removed user is refused; close(Selector::User) returns 1 |
a_speed_limit_change_keeps_the_users_sessions |
tests/hot_swap.rs |
The session ids before and after a limit change are the same |
removing_an_inbound_closes_its_sessions |
tests/hot_swap.rs |
removed names it; its connection closes; its port refuses; the other inbound is untouched |
a_tun_device_change_needs_allow_disruptive |
tests/hot_swap.rs |
A device change and a settings change are each refused as Disruptive without the flag. Needs CAP_NET_ADMIN; skipped (and passing) when a device cannot be created |
shutdown_ends_pending_grace_timers |
tests/hot_swap.rs |
One-hour removal and drain timers do not hold up shutdown(Duration::ZERO) |
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake |
tests/hot_swap.rs |
A client whose QUIC handshake never completes does not hold shutdown |
direct_flows_are_dialed_on_hooked_sockets |
tests/socket_policy.rs |
Direct TCP and UDP sockets pass the hook |
a_socks_upstream_is_dialed_on_hooked_sockets |
tests/socket_policy.rs |
A SOCKS upstream’s control connections and its UDP relay socket pass the hook |
a_refusing_hook_fails_the_dial |
tests/socket_policy.rs |
A hook error fails the SOCKS request |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
tests/tracking.rs |
Batches arrive about one interval apart; with the shutdown batch they hold every byte; take_usage is then empty and usage_snapshot still answers |
a_slow_usage_sink_does_not_stall_the_data_plane |
tests/tracking.rs |
A sink taking 1 s per report delays only its own batches |
cancelling_the_root_cancels_every_session |
tests/unit/session.rs |
Bound and unbound sessions both end when the root is cancelled |
a_stopped_listener_opens_no_session |
tests/unit/session.rs |
Sessions::open under a cancelled stop registers nothing |
the_empty_plane_drops_everything |
tests/unit/plane.rs |
Epoch 0; TCP and UDP, port 53 included, go to the blackhole |
the_first_plan_binds_and_builds_everything_then_publishes |
tests/unit/plan.rs |
The first plan has no disruptive step and ends with PublishPlane; every tag starts at version 1 |
each_rebuild_of_a_tag_takes_the_next_version, a_tag_added_back_takes_a_version_it_never_had |
tests/unit/plan.rs |
Version numbering |
no_tag_may_be_empty |
tests/unit/validate.rs |
The empty tag stays the supervisor’s |
Nothing asserts on the drop of the last handle (the socket-policy tests end that way without checking it), on ApplyError::Stopped, or on the command channel’s capacity. check is exercised through the app, for example: a_member_with_no_upstream_is_refused (app/tests/integration/e2e_balancer.rs) runs the binary’s --test on a balancer over a freedom outbound and expects a failure; a_backend_without_its_server_is_rejected (app/tests/integration/e2e_dns.rs) runs --test on a UDP DNS backend with no server and expects a failure; and every_lowered_node_builds_and_the_dns_setup_with_it (app/tests/unit/subscribe.rs) calls check on the lowered example subscribe file.
supervisor/benches/tracking.rs measures the cost of a usage take and the reconcile that feeds it over a server-sized ledger, and the per-byte overhead of a metered flow (cargo bench -p etemenanki-supervisor --bench tracking); see Tracking.