Skip to content

Listeners and the serve loop

Source files: 57 · checked against Etemenanki 555b7df · katana v4.1.1
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/system/mod.rs
  • Etemenanki/supervisor/src/system/listener.rs
  • Etemenanki/supervisor/src/topology/inbound/mod.rs
  • Etemenanki/supervisor/src/build/inbound.rs
  • Etemenanki/supervisor/src/lib.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/entity/session.rs
  • Etemenanki/supervisor/src/entity/usage.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/topology/spec_plan/inbound.rs
  • Etemenanki/supervisor/src/topology/spec_plan/plan.rs
  • Etemenanki/supervisor/src/build/apply.rs
  • Etemenanki/supervisor/src/build/users.rs
  • Etemenanki/supervisor/src/build/validate.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/tls/config.rs
  • Etemenanki/protocols/src/error.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/http/protocol.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/ss_2022/crypto.rs
  • Etemenanki/protocols/src/ss_aead.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/config.rs
  • Etemenanki/protocols/src/hysteria/server/endpoint.rs
  • Etemenanki/protocols/src/hysteria/obfs.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/tun/device.rs
  • Etemenanki/protocols/src/tun/config.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/app/src/main.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/supervisor/tests/unit/serve.rs
  • Etemenanki/supervisor/tests/unit/session.rs
  • Etemenanki/supervisor/tests/unit/plan.rs
  • Etemenanki/supervisor/tests/unit/validate.rs
  • Etemenanki/supervisor/tests/hot_swap.rs
  • Etemenanki/supervisor/tests/tracking.rs
  • Etemenanki/protocols/tests/pipeline/hysteria.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/tests/unit/trojan/core.rs
  • Etemenanki/app/tests/integration/e2e_unix.rs
  • Etemenanki/app/tests/integration/e2e_hysteria_inbound.rs
  • Etemenanki/app/tests/integration/e2e_tun.rs
  • Etemenanki/app/tests/integration/e2e_route_context.rs
  • katana/src/manager/node.rs

supervisor/src/system/listener.rs and supervisor/src/serve.rs are where an inbound’s spec becomes a bound socket and then running connections. A listener is bound while an apply is being prepared, so a port or device that cannot be bound refuses the whole spec before anything changes, and it starts serving only when the apply commits. It lives for as long as its BindSpec stays in the spec: a changed inbound on an unchanged bind gets a new handler on the same socket. Every accepted connection is admitted against two per-listener semaphores and registered as a session before its transport handshake starts, and is then driven through its protocol core by a small watchdog that bounds the handshake. Hysteria 2 and TUN listeners own a UDP socket or a device outright and hand the per-connection work to the protocols crate.

This page is for contributors who add an inbound protocol, change how connections are admitted, or touch anything that runs under a listener’s or a session’s token. Which listeners an apply binds, swaps or stops is decided on Planning and applying a change. How a session is bound to its user, revoked, closed and listed is on Users, principals and sessions. What the connector does with each flow is on The plane: routing each flow, and what a session’s byte counter feeds is on Per-user usage accounting.

Concern Where Notes
Bind in prepare system/listener.rs → bind, Pending::pair Opens a TCP or UDP socket, a Unix socket file or a TUN device, then builds what serving needs. Nothing is served yet.
Serve in commit system/listener.rs → Listener::start Spawns the task that serves a pending pair. It cannot fail.
Listener identity topology/inbound/mod.rs → BindSpec An unchanged bind keeps its socket and its admission counts across applies.
Handler swap Listener::prepare_swap, Listener::swap Stream listeners: the next socket accepted gets the new handler. A Hysteria 2 listener applies it inside its running endpoint, and a TUN listener restarts its runtime on it (see Handler swaps).
User tables Listener::prepare_users, Listener::store_users; build/inbound.rs → UserTable A user change stores a new table into the running handler; the handler is not rebuilt.
Accepting and admission serve.rs → run_stream_inbound Two semaphores per listener, taken without waiting.
Sessions at accept run_stream_inbound → Sessions::open Every accepted socket is a session before its transport handshake starts.
Transport fan-out serve.rs → Connection::serve_socket InboundTransport::accept yields one byte stream per socket, or one per HTTP/2 stream for gRPC.
Protocol dispatch serve.rs → serve_connection The SOCKS driver, or a ProxyServerRuntime over the protocol’s core.
Handshake watchdog serve.rs → drive HANDSHAKE_TIMEOUT per runtime step until the core is established.
Hysteria 2 serve.rs → run_hysteria_inbound One session per QUIC connection; the endpoint is changed in place.
TUN serve.rs → run_tun_inbound One runtime per device at a time, restarted on a swap; flows have no session.
Stopping Listener::stop Stops accepting. Stream and Hysteria 2 connections run on.

This code does not parse protocols (the cores in etemenanki-protocols do), route or dial (the AppConnector it hands each connection does), decide which listeners to bind, swap or stop (the plan does), or end sessions on its own (see Tokens and tasks for what does). Its only accounting is feeding each session’s Wire counter.

Every front end reaches it the same way, by applying a spec: etemenanki-app runs one supervisor for its config, and katana runs one per node. The serving code itself is not public. serve is a pub(crate) module in supervisor/src/lib.rs, system is public but its only module, listener, is pub(crate), and build::inbound is pub(crate). A front end sees the types in topology::inbound (BindSpec, TunSource, SuppliedTun and the handler types), but no public API takes a handler, so only the crate builds them (in build::inbound).

BindSpec: what is bound, and the listener’s identity

Section titled “BindSpec: what is bound, and the listener’s identity”
supervisor/src/topology/inbound/mod.rs
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BindSpec {
Tcp { host: String, port: u16 },
Udp { host: String, port: u16 },
Unix(PathBuf),
Tun(TunSource),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TunSource {
Create(DeviceSpec),
Fd(SuppliedTun),
}
#[derive(Debug, Clone)]
pub struct SuppliedTun {
pub fd: Arc<OwnedFd>,
pub mtu: u16,
}

The bind is the listener’s identity. plan matches each desired inbound to a running one by bind ==, not by tag, and Actor::prepare finds the running Listener the same way (listeners.iter().position(|l| l.bind == inbound.bind)). Two binds are the same when every field is equal. A DeviceSpec compares its name, mtu, addresses and routes, so any change to the device is a different bind. SuppliedTun implements PartialEq by hand: two are equal when they hold the same raw descriptor number and the same mtu. The code relies on a descriptor number being unique while it is open, and both are held open.

Validation makes the bind a usable key. validate refuses two inbounds with one bind (<bind> is bound by another inbound), and validate_inbound pairs each protocol with one kind of bind: SOCKS, HTTP, Trojan, VLESS, VMess, Shadowsocks and Shadowsocks 2022 with Tcp or Unix, Hysteria 2 with Udp, TUN with Tun (the protocol cannot be served on <bind>). A TCP and a UDP bind on the same port are different binds. The rules and their tests are on Validation.

BindSpec implements Display, which every listener log line and bind error uses:

Variant Displayed as Example
Tcp {host}:{port} 0.0.0.0:1080
Udp udp {host}:{port} udp 0.0.0.0:443
Unix unix:{path} unix:/run/etemenanki/socks.sock
Tun(Create) tun {name}, or tun auto when name is None tun ete0
Tun(Fd) tun fd {raw fd} tun fd 7

TunSource::Create is a device the inbound creates, addresses and routes itself, which needs CAP_NET_ADMIN on Linux; the kernel removes it when its last descriptor closes. TunSource::Fd is a device handed over already configured, such as a mobile VPN API’s descriptor, whose platform owns the interface, its addresses and its routes. The FFI crate builds this variant. The listener serves duplicates of fd and closes them when it stops; fd itself closes when its last holder, the spec included, drops it. TunSource::mtu() returns the MTU from either variant; validation requires at least 1280 (MIN_TUN_MTU).

supervisor/src/topology/inbound/mod.rs
pub enum InboundHandler {
Stream(Arc<StreamInbound>),
Hysteria2(Hy2Handler),
Tun(TunHandler),
}
pub struct StreamInbound {
pub tag: CompactString,
pub protocol: StreamProtocol,
pub transport: InboundTransport,
pub sniff: bool,
}
pub enum StreamProtocol {
Socks(ArcSwap<SocksInbound<Principal>>),
Http(ArcSwap<HttpServerConfig<Principal>>),
Trojan(ArcSwap<trojan::Validator<Principal>>),
Vless(ArcSwap<vless::Validator<Principal>>),
Vmess(Arc<AccountValidator<Principal>>),
Shadowsocks(ArcSwap<Resolved<Principal>>),
Ss2022(ArcSwap<Ss2022Users>),
}
pub struct Ss2022Users {
pub config: Arc<Ss2022ServerConfig<Principal>>,
pub validator: Option<Arc<ss_2022::Validator<Principal>>>,
}
pub struct Hy2Handler {
pub tag: CompactString,
pub connection: Arc<Hy2ConnectionConfig<Principal>>,
pub quic: quinn::ServerConfig,
pub obfs: Option<Obfs>,
pub max_connections: usize,
}
pub struct TunHandler {
pub tag: CompactString,
pub inbound: TunInbound<Principal>,
}

A handler is everything a listener serves with, split from the listener so that a changed inbound can get a new one on the same socket.

  • StreamInbound is what the accept loop serves every socket with. tag is what routes match on and what each session is filed under. transport is the layer under the protocol; a Unix listener carries none and never reads it. sniff says whether a core reads a domain out of a flow addressed by IP.
  • StreamProtocol holds what each connection’s core is built from. Every user-dependent part sits behind an ArcSwap, so a user change stores a new table into the running handler and never rebuilds it. VMess is the exception in form, not in effect: it keeps one Arc<AccountValidator> and swaps its user set itself (AccountValidator::set_users), which keeps its replay state. The type parameter is Principal, the per-user payload every protocol’s user table hands back on a successful handshake (see Users, principals and sessions). StreamProtocol::name() returns the name used in log lines: socks, http, trojan, vless, vmess, shadowsocks or shadowsocks-2022.
  • Ss2022Users is a Shadowsocks 2022 inbound’s user table: the server config with its key and users, and in multi-user mode the lookup for the identity header.
  • Hy2Handler is what a Hysteria 2 listener applies inside its one endpoint, because two QUIC endpoints cannot share a port. connection is the per-connection config, with the user table (authenticator) and the circuit budget (circuit_permits) inside it. quic is the TLS material new handshakes are answered with.
  • TunHandler is the TUN inbound itself, built whenever the inbound is new or changed.
supervisor/src/build/inbound.rs
pub(crate) fn build_handler<U: UserId>(
spec: &InboundSpec,
users: &Admissions<U>,
circuits: Option<&Arc<Semaphore>>,
) -> io::Result<InboundHandler>;
pub(crate) enum UserTable {
Socks(SocksInbound<Principal>),
Http(HttpServerConfig<Principal>),
Trojan(trojan::Validator<Principal>),
Vless(vless::Validator<Principal>),
Vmess(Vec<(Uuid, Arc<Principal>)>),
Shadowsocks(Resolved<Principal>),
Ss2022(Ss2022Users),
Hysteria2(Authenticator<Principal>),
None,
}
pub(crate) fn user_table<U: UserId>(spec: &InboundSpec, users: &Admissions<U>) -> io::Result<UserTable>;
impl UserTable {
pub(crate) fn fits(&self, protocol: &StreamProtocol) -> bool;
pub(crate) fn store_into(self, protocol: &StreamProtocol);
}

build_handler builds a whole handler when an inbound is new or changed. It builds the user table first, with user_table, and then:

Protocol What build_handler makes
Hysteria 2 hy2_handler: the masquerade (Masquerade::default(), or Masquerade::new(status, body, content_type)); circuit_permits, which is the running listener’s semaphore when circuits is Some and a new Semaphore::new(max_circuits) otherwise; the connection config (the authenticator in an ArcSwap, the masquerade, sniff, udp from udp_idle_timeout, the circuit permits); quic from hy2_endpoint::server_config(cert_pem, key_pem); obfs mapped from the spec; max_connections.
TUN TunInbound::new(TunConfig { user: Principal::anonymous(), mtu: device.mtu(), udp, udp_idle_timeout, max_flows }), with without_sniffing() unless the spec sniffs.
Every stream protocol A StreamInbound with the table in its StreamProtocol variant. The transport is InboundTransport::Tcp on a Unix bind and for the protocols that take no transport (SOCKS, Shadowsocks, Shadowsocks 2022). For HTTP, Trojan, VLESS and VMess, build_transport makes Tcp, Tls (no ALPN), Ws or Grpc. TLS comes from TlsServerConfig::from_pem(cert, key, alpn); when it is layered under WebSocket its ALPN is http/1.1, and under gRPC h2. InboundTransport::ws(path, host, tls) requires the request’s Host to match host when one is given.

A certificate or key that does not parse fails here and becomes ApplyError::Build, displayed as building inbound <tag> failed: <error>. etemenanki-app’s --test shows it for a Hysteria 2 certificate file with no certificate in it:

ERROR etemenanki_app: configuration invalid: building inbound hy2-in failed: hysteria2: the certificate file contains no certificates

hy2_endpoint::server_config builds a TLS 1.3-only rustls config with ALPN h3 and the QUIC transport parameters described on Hysteria 2: server. Its errors are hysteria2: could not read the certificate: <error>, hysteria2: the certificate file contains no certificates, hysteria2: could not read the private key: <error>, hysteria2: the key file contains no private key, hysteria2: certificate and key do not match: <error>, hysteria2 tls setup failed: <error> and hysteria2 quic tls setup failed: <error>.

Other checks in build_handler and user_table guard what validation, which plan runs before prepare, already refuses as ApplyError::Invalid (inbound <tag>: <reason>): the masquerade that Masquerade::new checks, each Shadowsocks 2022 user key that normalise_psk checks, and the combinations behind inbound <tag>: protocol and bind do not match and missing tls. etemenanki-app’s --test prints the validation text for the first two:

ERROR etemenanki_app: configuration invalid: inbound hy2-in: hysteria2: 233 is the authentication success status and cannot be used for the masquerade
ERROR etemenanki_app: configuration invalid: inbound ss-in: user alice@example.com: shadowsocks-2022: PSK too short (16 < 32)

A UserTable is the user-dependent part of a handler, built from the users the inbound admits. Which credential each protocol reads, what each table holds in an open or shared mode, and the labels user_table fills in are on Users, principals and sessions. fits says whether a table is the kind a StreamProtocol stores, and store_into stores it:

UserTable Stored by Into
Socks, Http, Trojan, Vless, Shadowsocks, Ss2022 store_into ArcSwap::store on the matching StreamProtocol cell
Vmess store_into AccountValidator::set_users on the running validator
Hysteria2 Listener::store_users Hy2Inbound::set_authenticator
None Listener::store_users Nothing: a TUN device admits no users
supervisor/src/system/listener.rs
pub(crate) enum Bound {
Stream(StreamListener),
Udp(std::net::UdpSocket),
Tun { device: OwnedFd, first: OwnedFd },
}
pub(crate) enum StreamListener {
Tcp(TcpListener),
Unix {
listener: tokio::net::UnixListener,
_file: SocketFile,
},
}
pub(crate) struct SocketFile {
path: PathBuf,
id: (u64, u64),
}
pub(crate) enum AcceptedSocket {
Tcp(tokio::net::TcpStream),
Unix(tokio::net::UnixStream),
}
impl StreamListener {
pub(crate) async fn accept(&self) -> io::Result<(AcceptedSocket, Option<IpAddr>)>;
}
  • StreamListener::accept returns the peer’s IP for TCP (peer.ip(); the port is dropped) and None for a Unix socket, which has no IP peer.
  • Dropping a StreamListener releases what it holds. So a listener bound by an apply that is then refused leaves nothing behind, and a stopped accept loop releases its port by dropping the listener.
  • SocketFile is the filesystem entry a Unix listener created, with the file’s device and inode. Its Drop removes the path only while symlink_metadata(path) still reports that device and inode, so a file someone else has since put at the path is left alone. A failed removal is logged at debug as could not remove <path>: <error> and otherwise ignored. _file is declared after listener, so the descriptor closes before the file is removed.
  • Bound::Tun carries the device’s descriptor, which the listener keeps, and first, a duplicate for the first runtime.
supervisor/src/system/listener.rs
pub(crate) enum Pending {
Stream(StreamListener, Arc<StreamInbound>),
Hysteria2 {
tag: CompactString,
inbound: Hy2Inbound<Principal>,
endpoint: OpenedEndpoint,
circuits: Arc<Semaphore>,
obfs: Option<Obfs>,
},
Tun {
device: OwnedFd,
prepared: PreparedDevice,
handler: TunHandler,
},
}
pub(crate) struct Listener {
pub tag: CompactString,
pub bind: BindSpec,
stop: CancellationToken,
kind: ListenerKind,
}
enum ListenerKind {
Stream(watch::Sender<Arc<StreamInbound>>),
Hysteria2 {
inbound: Hy2Inbound<Principal>,
tag: Arc<ArcSwap<CompactString>>,
circuits: Arc<Semaphore>,
obfs: Option<Obfs>,
},
Tun {
device: OwnedFd,
restart: mpsc::UnboundedSender<(TunHandler, PreparedDevice)>,
},
}
pub(crate) enum Swap {
Stream(Arc<StreamInbound>),
Hysteria2(Hy2Handler),
Tun(TunHandler, PreparedDevice),
}
pub(crate) struct Users(UserTable);
  • A Pending is a bound handle paired with the handler it will serve, with everything serving needs already built. It is the only thing Listener::start accepts, so starting cannot fail.
  • A Listener is a bound handle and the task serving it. tag is the tag of the inbound it currently serves, kept current by swap. bind is the key the actor finds it by. stop is a child of the supervisor’s root token, so shutdown reaches every listener.
  • ListenerKind holds what the actor needs to reach the running task: the watch sender of a stream listener’s handler; the Hysteria 2 inbound handle (a clone of the one the task runs), a shared cell with its current tag, its circuit budget and its obfuscation; the TUN device’s own descriptor and the sender of the runtime’s restart channel.
  • Swap and Users are a handler and a user table already checked against the listener they go to. They are made in prepare and consumed in commit.
supervisor/src/serve.rs
#[derive(Clone)]
pub(crate) struct Shared {
pub plane: PlaneCell,
pub sessions: Arc<Sessions>,
pub flows: Tracker,
pub tracker: TaskTracker,
}
impl Shared {
fn connector(&self, tag: CompactString, source: Option<IpAddr>, wire: Wire, stop: &CancellationToken)
-> Option<(AppConnector, Arc<Session>)>;
fn connector_for(&self, tag: CompactString, source: Option<IpAddr>, session: Arc<Session>)
-> AppConnector;
fn spawn_until<F>(&self, token: CancellationToken, fut: F)
where
F: Future<Output = ()> + Send + 'static;
}
struct Connection {
inbound: Arc<StreamInbound>,
peer: Option<IpAddr>,
carrier: Arc<Session>,
stop: CancellationToken,
shared: Shared,
live: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
}
struct WireCounted<S> {
inner: S,
wire: Option<Arc<Wire>>,
}

Shared is what every accept loop shares with the supervisor, cloned into each task. Its fields:

Field Type Used for
plane PlaneCell The routes and outbounds every connector dials through (The plane)
sessions Arc<Sessions> The session registry: open at accept, and root() for Hysteria 2 shutdown
flows Tracker (crate::track) Where every flow is registered and metered (Tracking)
tracker tokio_util::task::TaskTracker Every accept loop and connection task, so shutdown can wait for them

The two trackers are unrelated: flows is the supervisor’s flow registry, and tracker is Tokio’s task tracker. Its methods:

Method What it does
connector Registers a new session with Sessions::open and returns a connector for its flows, or None when the listener’s stop is already cancelled.
connector_for Builds a connector for a session already registered: AppConnector::new(plane, FlowContext { inbound_tag, source }, flows, Some(session)).
spawn_until Spawns fut on the task tracker inside a tokio::select! against token.cancelled(). When the token fires first, the future is dropped where it stands.

Connection is one accepted socket and what serving it needs: the handler, the peer IP, the socket’s own session (carrier), the listener’s token, Shared, and the two admission permits.

WireCounted wraps the SOCKS driver’s stream. poll_read adds the bytes a successful read filled as up, poll_write adds the count a successful write returned as down, and flush and shutdown pass through. What that counts is on Per-user usage accounting.

sequenceDiagram
  participant A as Actor (apply)
  participant B as listener bind
  participant P as Pending pair
  participant L as Listener start
  participant T as serving task
  Note over A: prepare, fallible, nothing visible
  A->>B: bind(tag, BindSpec) for each new bind
  B-->>A: Bound socket, Unix file or TUN descriptors
  A->>P: pair(bound, handler)
  P-->>A: Pending with Hysteria 2 endpoint opened or TUN device prepared
  Note over A: a later failure drops every Pending, which releases it
  Note over A: commit, infallible
  A->>L: start(bind, pending, shared, root child token)
  L->>T: spawn run_stream_inbound, run_hysteria_inbound or run_tun_inbound
  L-->>A: Listener, logged as inbound tag listening on bind
flowchart TB
  root["root token (Sessions::root)"]
  stop["listener stop token, child of root"]
  sess["session token, child of root"]
  acc["run_stream_inbound"]
  sock["serve_socket, one per accepted socket"]
  conn["serve_connection, one per byte stream"]
  hy["run_hysteria_inbound"]
  hyconn["QUIC connection tasks (JoinSet in Hy2Inbound::run)"]
  tun["run_tun_inbound"]
  tunflow["flow tasks (JoinSet in TunInbound::run)"]
  root --> stop
  root --> sess
  stop -. ends .-> acc
  stop -. drains .-> hy
  stop -. ends .-> tun
  acc --> sock
  sock --> conn
  sess -. drops .-> sock
  sess -. drops .-> conn
  sess -. closes .-> hyconn
  hy --> hyconn
  tun --> tunflow

The accept loops, the Hysteria 2 and TUN listener tasks and every stream connection task run on the supervisor’s TaskTracker. The per-connection tasks of a Hysteria 2 listener and the per-flow tasks of a TUN device live in a JoinSet owned by the protocols crate’s run future, so they end when it does. A listener’s token stops only its accepting task. A connection runs under its session’s token, which is a child of the root and not of the listener, so stopping a listener leaves the connections it accepted running. Nothing in the serving code cancels a session’s token. It is cancelled by a removal policy, Supervisor::close, the removal of its inbound, shutdown, or a first flow that presents a principal already revoked (see Users, principals and sessions); the session’s own Drop cancels it too. TUN is the exception: its flows are tasks of the runtime the listener’s token stops.

The sequence follows one client of a VLESS-over-TLS inbound from accept to the end of its relay.

sequenceDiagram
  participant C as Client
  participant L as run_stream_inbound
  participant S as Sessions
  participant K as serve_socket task
  participant T as InboundTransport
  participant V as serve_connection task
  participant D as drive and runtime
  C->>L: TCP connect
  L->>L: try_acquire live permit, then handshake permit
  L->>S: open(tag, peer, counted wire, stop)
  S-->>L: carrier session, token is a child of root
  L->>K: spawn_until(session token, serve_socket)
  K->>T: accept(tcp, sink)
  T->>C: TLS handshake within TRANSPORT_HANDSHAKE_TIMEOUT
  T->>V: sink(stream) spawns serve_connection on the same session
  V->>D: ProxyServerRuntime over VlessCore
  C->>D: request header
  Note over D: first flow admitted on the session binds its user
  Note over D: core established, drive drops its handshake permit
  D->>D: relay, adding transport bytes to the session wire
  Note over V: the task ends, dropping its permits and the session
supervisor/src/system/listener.rs
pub(crate) async fn bind(tag: &str, bind: &BindSpec) -> io::Result<Bound>;

Actor::prepare calls bind for every desired inbound that has no running listener on its bind, and maps an error to ApplyError::Bind { inbound, bind, source }, displayed as inbound <tag>: binding <bind> failed: <error>. tag is used only in log lines.

BindSpec What bind does Returns
Tcp { host, port } std::net::TcpListener::bind((host, port)), set_nonblocking(true), tokio::net::TcpListener::from_std Bound::Stream(StreamListener::Tcp)
Udp { host, port } std::net::UdpSocket::bind((host, port)) Bound::Udp
Unix(path) See below Bound::Stream(StreamListener::Unix)
Tun(Create(spec)) etemenanki_protocols::tun::open(spec).await, then try_clone for first. Logs inbound <tag> owns tun device <name> at info. Bound::Tun
Tun(Fd(supplied)) etemenanki_protocols::tun::adopt(supplied.fd.as_fd()), a non-blocking duplicate; then try_clone for first. Nothing is created, addressed or routed, so no privilege is needed. Logs inbound <tag> serves a supplied tun device at info. Bound::Tun

host does not have to be an IP literal. Both bind calls take (host, port) through ToSocketAddrs, which resolves a name with the system resolver, and bind the first resolved address that succeeds; the error is the last address’s when none does.

tun::open (protocols/src/tun/device.rs) does the device work:

  1. check_platform(spec): off Linux, a spec with routes is refused with tun routes are installed only on Linux; add them with the OS route tool. Validation calls the same function for a TunSource::Create, so --test and a start agree.
  2. Build the device with the MTU, the name when one is given, the first IPv4 address and every IPv6 address (ipv6_tuple); each further IPv4 address is added to the live device.
  3. Set it non-blocking and read back the name and interface index the kernel settled on.
  4. On Linux, install each route through rtnetlink as a link-scope route on that interface. Routes are never removed explicitly: they belong to the interface, and the kernel drops them with it when its last descriptor closes.

Any failure drops the half-built device, which destroys it. A duplicate shares its original’s file status flags, so adopt makes the supplied descriptor non-blocking too. On Android and iOS, tun::open fails with tun devices are created by the system VPN API here; adopt its descriptor. The device side is on TUN.

  1. symlink_metadata(path) looks at what is at the path, without following a symlink:
    • a socket is taken to be stale, left by a crashed run, and is removed. A live listener of the same supervisor on the same path has the same bind and is kept, so bind never meets its own socket;
    • anything else, a symlink included, is someone else’s file: bind fails with io::ErrorKind::AlreadyExists and <path> exists and is not a socket;
    • NotFound continues; any other error is returned.
  2. std::os::unix::net::UnixListener::bind(path) creates the socket file.
  3. symlink_metadata(path) again records the new file’s device and inode in a SocketFile. If that fails, the file is removed and the error returned.
  4. set_nonblocking(true) and tokio::net::UnixListener::from_std. A failure here drops the SocketFile, which removes the file.

etemenanki-app shows both outcomes. With listen set to a path that holds a regular file, --test passes, because a dry run binds nothing, and a start fails:

ERROR etemenanki_app: failed to start: inbound in: binding unix:/run/etemenanki/socks.sock failed: /run/etemenanki/socks.sock exists and is not a socket

With a free path, the start logs the listener and removes the file again on SIGTERM:

INFO etemenanki_supervisor::system::listener: inbound in listening on unix:/run/etemenanki/socks.sock

A Unix listener has no local IP and no peer IP. It carries no transport: validation refuses any shape other than plain TCP on it (a unix socket carries no transport; its shape must be plain tcp), and SOCKS with udp on it needs udp_bind (socks over a unix socket has no local IP for UDP associate; set udp_bind or turn udp off). What a UDP ASSOCIATE over a Unix socket may hear is on SOCKS.

supervisor::check(spec) runs Actor::prepare with bind = false on a fresh actor: every handler, user table, outbound and route is built, but no new listener is bound or paired, so no port is taken, no Hysteria 2 endpoint is opened and no TUN device is created. It refuses what a start would, short of a port or device that cannot be bound. etemenanki-app’s --test runs it and prints Configuration OK. on success, or logs configuration invalid: <error> at error and exits with a failure status.

supervisor/src/system/listener.rs
impl Pending {
pub(crate) fn pair(bound: Bound, handler: InboundHandler) -> io::Result<Self>;
}
Bound + InboundHandler What pair builds
Stream + Stream Nothing more: Pending::Stream(listener, inbound)
Udp + Hysteria2 check_obfs (a Salamander key Salamander::new would refuse fails here, not at commit); circuits is connection.circuit_permits; Hy2Inbound::new(ListenerConfig { connection, quic, obfs, max_connections }); inbound.open(socket) builds the QUIC endpoint on the socket.
Tun + Tun handler.inbound.prepare(first): registers the duplicate with the Tokio runtime and configures the IP stack (MTU, UDP timeout, packet information on macOS and iOS).
Any other pair io::ErrorKind::InvalidInput, the handler is not of the kind its listener serves

An opened endpoint admits nobody until it is served. Hy2Inbound::open:

  1. stores the Salamander key of the inbound’s obfs, or none, into the inbound’s shared SalamanderSwitch, so the endpoint obfuscates from its first packet. open fails on a key Salamander::new refuses (shorter than MIN_PSK_LEN, 4 bytes), although pair has already refused one with check_obfs;
  2. calls endpoint::from_socket(socket, None, switch, gate) with gate = RecvGate::default(), a closed gate. from_socket sets the socket non-blocking, wraps it for quinn’s Tokio runtime, and wraps it again in SalamanderSocket::switchable even when there is no key, so obfuscation can be switched on, off or to another key under the running endpoint. With no server config the endpoint refuses handshakes, and with the gate closed it reads nothing off the socket, so clients’ packets wait in the kernel as they would for a socket nobody serves yet. run opens the gate with RecvGate::open.

A PreparedDevice is likewise read by nothing until its runtime starts. TunInbound::prepare registers the descriptor with the reactor through TunDevice::new and sets the stack’s IpStackConfig. Its mtu setter refuses less than 1280, which is why validation’s MIN_TUN_MTU is 1280; prepare turns that refusal into InvalidInput. Dropped unserved, an opened endpoint and a prepared device both close, and the inbound is as it was.

The mismatch cannot happen after validation, which pairs every protocol with its kind of bind; the check keeps start infallible whatever the caller passes. The error, like every failure of pair (a refused key, from_socket or quinn’s endpoint creation failing, TunInbound::prepare failing), becomes ApplyError::Build.

supervisor/src/system/listener.rs
impl Listener {
pub(crate) fn start(bind: BindSpec, pending: Pending, shared: &Shared, stop: CancellationToken) -> Self;
}

Actor::commit calls it for every pending pair with stop = root.child_token(), after the plane is published and removed users are revoked, and appends the listener to Actor::listeners.

Pending start spawns on the task tracker ListenerKind
Stream run_stream_inbound(listener, rx, shared, stop), with rx the receiving end of watch::channel(inbound) Stream(tx)
Hysteria2 run_hysteria_inbound(inbound.clone(), tag_cell, endpoint, shared, stop), with tag_cell = Arc::new(ArcSwap::from_pointee(tag)) Hysteria2 { inbound, tag: tag_cell, circuits, obfs }
Tun run_tun_inbound(handler, prepared, rx, plane, flows, stop) through spawn_tun, with rx the receiving end of an unbounded mpsc channel Tun { device, restart: tx }

Every start logs inbound <tag> listening on <bind> at info, target etemenanki_supervisor::system::listener.

A desired inbound whose bind is running either is unchanged or gets a new handler (Step::SwapHandler), and an unchanged one may still get a new user table. The plan’s rules are on Planning and applying a change. The listener side:

supervisor/src/system/listener.rs
impl Listener {
pub(crate) fn prepare_swap(&self, handler: InboundHandler) -> io::Result<Swap>;
pub(crate) fn swap(&mut self, swap: Swap, shared: &Shared);
pub(crate) fn prepare_users(&self, table: UserTable) -> io::Result<Users>;
pub(crate) fn store_users(&self, Users(table): Users);
pub(crate) fn circuits(&self) -> Option<&Arc<Semaphore>>;
}
Listener prepare_swap (prepare) swap (commit)
Stream Refuses a handler of another kind Sets tag, then send_replace(inbound) on the watch sender. The accept loop reads the handler at each accept, so the next socket gets the new one.
Hysteria 2 Refuses another kind; check_obfs on the new key If obfs differs: set_obfs (its expect states that prepare_swap checked the key). Then tag and the tag cell, the recorded circuit budget, set_connection_config, set_quic and set_max_connections.
TUN Refuses another kind; duplicates the listener’s device descriptor and prepares the duplicate Sets tag and sends the handler and prepared duplicate on restart. If the send fails because the runtime task has already ended, spawn_tun starts a new runtime task on them and replaces the sender.

a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections pins a stream swap with a SOCKS inbound changed to HTTP on the same port: a relay opened before the apply keeps relaying, a new client speaks HTTP, and a socket accepted just before the apply but handshaking after it is still answered as SOCKS.

What a Hysteria 2 swap reaches is set by each setter of Hy2Inbound:

Setter Reaches
set_connection_config Connections that arrive afterwards.
set_quic New handshakes: the running endpoint takes the server config at once. Live connections keep what they negotiated.
set_max_connections Arrivals from now on. Lowering it closes no live connection; new ones are refused until enough have ended. The limit is a count checked per arrival, not a semaphore, so it can change under live connections.
set_obfs Every packet, both ways, from the moment it returns. Every connection negotiated under the old key is closed, since it cannot continue under the new one; the endpoint stays bound and keeps accepting. Setting the key already in use closes nothing.

The circuit budget is carried over when the running and the desired inbound on the bind are both Hysteria 2 with the same max_circuits (same_circuits in supervisor.rs): Actor::prepare passes Listener::circuits() to build_handler, so the new connection config carries the running semaphore and live circuits keep counting against it. A different size gets a new semaphore, a fresh budget beside the circuits still held under the old one.

A changed Hysteria 2 obfuscation, and a change to a TUN inbound’s device or settings, are Step::Disrupt, because they end live connections. Such a step needs ApplyOptions::allow_disruptive: an apply without it is refused with ApplyError::Disruptive.

etemenanki-app refuses such a reload and logs reload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply it; a restart applies it. See etemenanki-app: running, reloading and shutting down.

katana applies its specs with allow_disruptive set, so an obfuscation change applies there; katana has no TUN inbound.

A masquerade, certificate, UDP or connection-limit change on a Hysteria 2 inbound is a plain swap, not a Step::Disrupt (a_hysteria2_masquerade_change_swaps_without_disrupting).

Listener prepare_users accepts store_users
Stream A table that fits the handler’s protocol store_into the handler’s protocol
Hysteria 2 UserTable::Hysteria2 Hy2Inbound::set_authenticator(Arc::new(authenticator)), stored into the current connection config: connections authenticated from now on use it
TUN UserTable::None Nothing

Any other table is refused with the handler is not of the kind its listener serves. Both the apply path and the user-edit path (set_users, upsert_user, remove_user) end here. Neither stores a table before every table of the change is prepared, and both revoke removed users’ principals first, so a handshake against a table about to be replaced finds its principal revoked. That ordering and what it does to live sessions are on Users, principals and sessions.

In one apply, an inbound gets a swap or a user table, never both: a changed inbound’s handler is built with the new users already, and only an unchanged one whose admissions differ gets a table. So set_connection_config and set_authenticator never race on one Hysteria 2 listener, and since the actor runs one command at a time, no other caller can race them either.

supervisor/src/serve.rs
const ACCEPT_ERROR_BACKOFF: Duration = Duration::from_millis(100);
const MAX_HANDSHAKES_PER_INBOUND: usize = 204_800;
const MAX_LIVE_CONNECTIONS_PER_INBOUND: usize = 4_194_304;
pub(crate) async fn run_stream_inbound(
listener: StreamListener,
handler: watch::Receiver<Arc<StreamInbound>>,
shared: Shared,
stop: CancellationToken,
);

The loop is a biased tokio::select! of stop.cancelled() over listener.accept(). Biased means the stop is checked first on every turn, so once the listener is stopped not one more socket is accepted. On stop the loop breaks and drops the listener, which closes the port or removes the socket file.

supervisor/src/serve.rs
fn should_backoff_accept_error(e: &io::Error) -> bool;
async fn backoff_or_cancelled(token: &CancellationToken) -> bool;
io::ErrorKind Meaning Handling
ConnectionAborted, Interrupted One client went away before accept, or a signal interrupted the call. The listener is fine. debug accept error: <error>, then accept again at once
Anything else Usually resource exhaustion, such as running out of file descriptors, which would fail again at once warn accept error, backing off 100ms: <error>, then backoff_or_cancelled sleeps ACCEPT_ERROR_BACKOFF

The backoff keeps a persistent error from turning the loop into a busy spin that floods the log. backoff_or_cancelled races the sleep against the stop token and returns true when the token wins, and the loop then breaks without finishing the sleep. An accept error never ends the loop by itself.

Both semaphores are created when run_stream_inbound starts, so their counts belong to the listener and last as long as it does. A handler swap keeps them, since the loop keeps running.

Semaphore Capacity Taken Released
live MAX_LIVE_CONNECTIONS_PER_INBOUND (4,194,304) First, at accept When the connection has ended
handshakes MAX_HANDSHAKES_PER_INBOUND (204,800) Second, at accept For a protocol drive serves, drive drops its handle to the permit once its core reports established

Both are taken with try_acquire_owned: the loop never waits for a permit. When one is exhausted the loop logs at debug and continues, which drops the freshly accepted socket (the client sees the connection close) and any permit taken for it:

dropping inbound connection; live connection limit reached
dropping inbound connection; handshake limit reached

Refusing at once keeps the loop responsive and keeps a burst from piling up sockets nobody serves. The live cap is a guardrail, not a quota: a busy inbound carries hundreds of thousands of mostly idle connections, so the cap sits far above normal traffic. It is also above typical per-process file-descriptor limits, so do not expect it to act before descriptors run out. Running out shows up as an accept error, which the 100 ms backoff above handles. Operator-facing limits are on Limits and timeouts.

With both permits taken, the loop:

  1. reads the current handler, handler.borrow().clone(), which holds the watch channel’s read lock only while the Arc is cloned;
  2. opens the socket’s session, shared.sessions.open(inbound.tag, peer, Wire::counted(), &stop). The session is registered, and listed by Supervisor::sessions, from this moment, with no user until its first flow binds one. Its transport handshake is closed with its inbound too. open checks the stop token under the registry’s lock and returns None once it is cancelled; the loop then breaks, dropping the socket. That check is what lets a removed inbound’s session sweep miss nothing (a_stopped_listener_opens_no_session);
  3. builds a Connection with the handler, the peer, the session as carrier, clones of stop and Shared, and both permits;
  4. spawns connection.serve_socket(socket) with shared.spawn_until(carrier.token(), …).

The peer IP is carried rather than dropped at accept because routing needs it: RouteMatch::SourceCidr has nothing to match without it, and it is the one client a SOCKS UDP ASSOCIATE hears. The inbound tag feeds RouteMatch::InboundTag the same way. Both reach the route through FlowContext (The plane).

serve_socket runs the transport over the accepted socket and serves every byte stream it yields.

  • AcceptedSocket::Unix: serve_connection runs inline on the Unix stream, with local_ip = None and a connector for the carrier session with source = None. No transport runs and no keepalive is set.
  • AcceptedSocket::Tcp: local_ip is taken from tcp.local_addr(). The SOCKS driver puts it in its replies and, unless udp_bind is set, binds the UDP ASSOCIATE relay on it. Then InboundTransport::accept(tcp, sink) runs.
protocols/src/transports/accept.rs
pub const TRANSPORT_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
impl InboundTransport {
pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);
}

accept enables TCP keepalive on the socket, runs the transport’s own handshake within TRANSPORT_HANDSHAKE_TIMEOUT, and calls sink for each byte stream:

InboundTransport Streams
Tcp One, the socket itself
Tls One, after the TLS handshake
Ws One, after the WebSocket upgrade (over TLS when configured)
Grpc One per HTTP/2 stream on the tunnel service’s paths

A timeout fails with io::ErrorKind::TimedOut and tls handshake timed out, websocket handshake timed out or grpc handshake timed out. Any transport error ends serve_socket with debug inbound transport failed: <error>, and the socket and its session are dropped. The transports are on TCP and TLS transports and WebSocket and gRPC transports.

The sink closure decides each stream’s session:

  • One stream per socket (TCP, TLS, WebSocket). The stream is served as the carrier session: connector_for(tag, peer, carrier), and serve_connection is spawned under the carrier’s token.
  • gRPC. Each HTTP/2 stream is a session of its own, opened through Shared::connector with a Wire of its own.
supervisor/src/serve.rs
async fn serve_connection<S>(
inbound: Arc<StreamInbound>,
stream: S,
local_ip: Option<IpAddr>,
connector: AppConnector,
live: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
) where
S: AsyncRead + AsyncWrite + Unpin + Send + 'static;

serve_connection reads source from the connector (the peer IP, or None over a Unix socket) and sniff from the handler, and matches on the protocol:

StreamProtocol Driver Core Runtime BUF
Socks SocksInbound::serve on socks.load_full(), over WireCounted None n/a
Http drive HttpCore::new(config.load_full(), sniff, source) HttpCore::<Principal>::BUF_SIZE = MAX_HEAD = 65,536
Trojan drive TrojanCore::new(validator.load_full(), sniff, source) 16,384
Vless drive VlessCore::new(validator.load_full(), sniff, source) 16,384
Vmess drive VMessCore::new(validator.clone(), now_unix, sniff, source) 32,768
Shadowsocks drive ShadowsocksCore::new(resolved.load_full(), sniff, source) 20,480
Ss2022 drive Ss2022Core::with_system_clock(users.config, users.validator, sniff, source), after users.load_full() MAX_RECORD_LEN = 65,569

Each core is instantiated on Principal, and each BUF_SIZE is an associated constant of the core, passed through as the runtime’s const generic so each protocol gets buffers sized for its largest frame. Shadowsocks 2022’s is the largest record a peer may send: MAX_PACKET_SIZE (65,535) plus RECORD_OVERHEAD (a 2-byte length and two 16-byte tags, 34 bytes). The runtime allocates three buffers of BUF bytes per connection (transport read, transport staging, outbound scratch), so these constants set most of a connection’s fixed memory: 196,608 bytes for HTTP, 196,707 for Shadowsocks 2022, 98,304 for VMess, 61,440 for Shadowsocks, and 49,152 for Trojan and VLESS. The buffers are described on The server runtime.

Whatever the driver returns, an error is logged at debug with the protocol name and the client address, and an Ok is silent:

vless connection from Some(203.0.113.7) ended: inbound handshake timed out after 10s
socks connection from None ended: client did not complete its request in time

SOCKS does not fit a core: its UDP side lives on a second socket that the control connection only keeps alive. It gets its own driver:

protocols/src/socks/server.rs
impl<T: Send + Sync + 'static> SocksInbound<T> {
pub async fn serve<S, C>(
&self,
stream: S,
local_ip: Option<IpAddr>,
source: Option<IpAddr>,
connector: C,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
C: Connector<Flow<T>>,
C::Datagram: DatagramLink<Addr = Destination>;
}
  • The driver bounds its whole handshake with its own tokio::time::timeout(HANDSHAKE_TIMEOUT, …) and fails with client did not complete its request in time. It then relays a CONNECT or drives a UDP ASSOCIATE; see SOCKS.
  • No runtime reports what the stream moved, so serve_connection wraps it in WireCounted with the session’s wire, and the SOCKS greeting, authentication and request count toward the session like any other bytes.
  • The connector is passed as connector.charging_datagrams(). A UDP association’s datagrams never cross the control stream, so the payload of each datagram its sub-links send or receive is added to the session’s wire instead (Per-user usage accounting).
supervisor/src/serve.rs
pub trait Established {
fn is_established(&self) -> bool;
}
type Step = Result<Traffic, RuntimeError<io::Error>>;
async fn drive<const BUF: usize, Core, S>(
stream: S,
core: Core,
connector: AppConnector,
handshake: Arc<OwnedSemaphorePermit>,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
Core: ProxyCoreDecode<Target = Flow, Error = io::Error, TransportAddr = ()> + Established;

drive builds ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress(), a Stream that yields one Traffic delta per unit of work. It needs those yields: between two of them it can ask runtime.core().is_established() and stop watching at the right moment.

stateDiagram-v2
  [*] --> Handshake
  Handshake --> Handshake: Some(Ok), counted, core not established
  Handshake --> Relay: core established, handshake permit dropped
  Handshake --> TimedOut: no step within HANDSHAKE_TIMEOUT
  Handshake --> Failed: Some(Err)
  Handshake --> Ended: None
  Relay --> Relay: Some(Ok), counted
  Relay --> Failed: Some(Err)
  Relay --> Ended: None
  TimedOut --> [*]
  Failed --> [*]
  Ended --> [*]
  1. Handshake. While the core is not established, each runtime.next() is wrapped in tokio::time::timeout(HANDSHAKE_TIMEOUT, …).
    • Some(Ok(traffic)) is counted and the loop goes on.
    • A timeout returns io::ErrorKind::TimedOut with inbound handshake timed out after 10s (the text is built from HANDSHAKE_TIMEOUT.as_secs()).
    • Some(Err(e)) returns io::Error::other(e.to_string()), so the log shows the runtime’s text. RuntimeError prefixes a transport failure with transport: and a core error with proxy core: (for example proxy core: client did not complete its request in time); its other variants are the runtime’s own texts: effect targets an unknown outbound key, open reuses a live outbound key, forward range outside the event slice or held buffer, core consumed an impossible byte count, protocol frame exceeds the read buffer, stream effect on a datagram outbound or vice versa and bytes staged toward a datagram transport without a peer. When each is raised is on The server runtime.
    • None means the runtime finished before establishing, for example because the client closed or the core refused the request and finished. drive returns Ok(()).
  2. Relay. Once the core is established, drive drops its handshake permit (handshake.take() on the Option it holds it in) and polls the runtime until it ends, counting every step and converting a RuntimeError the same way.

Counting feeds the session: each step’s transport_rx is added to the session’s Wire as up and transport_tx as down, in both phases. Those are the bytes of the protocol stream after its transport, handshake included. What they mean for billing is on Per-user usage accounting.

Every stream core answers is_established from its Timing (protocols/src/core/mod.rs): true in Phase::Relay and Phase::Closing. A core enters Phase::Relay once its request is parsed and its flow is opening, so established means the request is accepted, not that the outbound is connected. A core still collecting a sniffing prefix (Phase::Sniff, bounded by SNIFF_TIMEOUT, 300 ms) is not established, so the watchdog covers sniffing too.

A runtime delivers no event to its core until the client sends something, and a core arms its handshake deadline at its first byte event. A client that connects and says nothing would therefore sit in the runtime forever. drive closes that gap from outside, and the two timers complement each other:

Timer Armed by Measures Catches
drive’s tokio::time::timeout Every runtime.next() before establishment Time between runtime steps A client that never speaks, or goes quiet
The core’s Timing deadline The first byte event, once Time since the request started A client that keeps sending but never completes its request

Both use HANDSHAKE_TIMEOUT from protocols/src/core/mod.rs, 10 s.

is_established is an inherent method on each core, and a generic function cannot call an inherent method. serve.rs declares Established and implements it with the established! macro, which forwards to the inherent method for HttpCore<Principal>, TrojanCore<Principal>, VlessCore<Principal>, VMessCore<Principal>, ShadowsocksCore<Principal> and Ss2022Core<Principal>. A new core drive should serve must be added to that list.

A Hysteria 2 inbound owns its UDP socket and QUIC endpoint for as long as its bind is in the spec. Everything inside a QUIC connection, from the credential exchange to the proxy streams and datagrams, is the protocols crate’s; see Hysteria 2: server. This section covers what the supervisor does around it.

supervisor/src/serve.rs
pub(crate) async fn run_hysteria_inbound(
inbound: Hy2Inbound<Principal>,
tag: Arc<ArcSwap<CompactString>>,
endpoint: OpenedEndpoint,
shared: Shared,
stop: CancellationToken,
);
protocols/src/hysteria/server/inbound.rs
impl<T> Hy2Inbound<T> {
pub fn open(&self, socket: std::net::UdpSocket) -> io::Result<OpenedEndpoint>;
pub async fn shutdown(&self);
}
impl<T: Send + Sync + 'static> Hy2Inbound<T> {
pub async fn run<C, F>(&self, endpoint: OpenedEndpoint, make_connector: F, token: CancellationToken)
where
F: Fn(IpAddr, ConnectionBytes) -> (C, CancellationToken) + Send + Sync + 'static,
C: Connector<Flow<T>> + Clone + Send + 'static,
C::Future: Send,
C::Stream: Send,
C::Datagram: DatagramLink<Addr = Destination> + Send;
}
#[derive(Clone)]
pub struct ConnectionBytes(quinn::Connection);
impl ConnectionBytes {
pub fn bytes(&self) -> (u64, u64);
}

In prepare, build_handler builds the QUIC server config from the certificate and key, listener::bind binds the UDP socket, and Pending::pair opens the endpoint on it, so a failure in any of them refuses the apply before anything changes. Once the task started at commit runs, run installs the QUIC server config and records the endpoint under the inbound’s lock. So a set_quic called before that is the config run installs, and one called after it is applied to the running endpoint. run then opens the receive gate and loops over a biased tokio::select!, in this order:

  1. the token, once: cancelling it switches run to draining;
  2. a finished connection task, so finished tasks are reaped before new clients pile on;
  3. the next incoming client. endpoint.accept() yielding None means the endpoint has closed, and the loop ends.

Each client takes a place under max_connections through an Admitted guard: a fetch_update on a shared count that succeeds only while fewer than the limit, read at that arrival, are taken. The guard moves into the connection’s task and gives the place back when the task ends. A client that finds no place, or arrives while run is draining, is refused with a QUIC refusal. An admitted connection’s task also records the obfuscation token current at its arrival, which is how set_obfs closes exactly the connections negotiated under the key it replaces. A draining run returns once its last connection task has ended; one whose endpoint closed first waits out the remaining tasks and returns.

run calls make_connector once per connection, as soon as its QUIC handshake completes, with the client’s IP and a ConnectionBytes. The supervisor’s closure:

  1. reads the inbound’s current tag from the tag cell. A listener kept across an apply may serve a renamed inbound, so the tag is read per connection;
  2. opens a session with Shared::connector(tag, Some(ip), Wire::quic(bytes), &stop) and returns the connector and the session’s token. The connector is cloned for each of the connection’s proxy streams and for its datagram channel, so they are all one session;
  3. when the listener was stopped while the client was in its handshake, returns a connector with no session and a token that is already cancelled. The connection is closed at once, so no session escapes the close of a removed inbound.

Cancelling the returned token closes that one connection with QUIC application error code 0x107 (DISCONNECT_CODE, HTTP/3’s H3_EXCESSIVE_LOAD) and the reason disconnected by the server, and nothing else. A connection closed because the obfuscation changed gets the same code with obfuscation changed.

ConnectionBytes holds the quinn::Connection and reads (udp_rx.bytes, udp_tx.bytes) from its stats: every UDP datagram of the connection, handshake, QUIC framing and TLS included. Wire::quic reads it whenever the session’s bytes are read, so a Hysteria 2 session counts its QUIC connection’s bytes rather than the payload it carried (a_hysteria2_session_bills_its_quic_connections_bytes). Because it holds the connection, it keeps reading the connection’s counts for as long as the session’s wire holds it.

run_hysteria_inbound runs inbound.run(endpoint, make, stop) and races it against the session registry’s root token:

  • run returns first. Cancelling stop makes run drain: clients that arrive afterwards are refused, and the connections already served carry on until they end or are closed through their session. When the last one has ended, run returns, and inbound.shutdown() releases the port.
  • The root is cancelled first. That happens only at Supervisor::shutdown, once its grace is over. A client still in its QUIC handshake has no session a root cancel could reach, so the task runs run and inbound.shutdown() together: shutdown closes the endpoint under the draining run, which ends every connection and handshake on it at once, and run returns with them. shutdown_does_not_wait_out_a_stalled_hysteria2_handshake pins that a stalled handshake does not hold shutdown up.

Hy2Inbound::shutdown closes the endpoint with application code 0x100 (CLOSE_CODE) and waits, bounded, for its port to come free; it returns at once when no endpoint is running. Its steps and timeouts are on Hysteria 2: server.

A TUN inbound owns a layer-3 device. A userspace IP stack in the protocols crate terminates every TCP connection and UDP flow the OS steers into it; see TUN.

supervisor/src/serve.rs
pub(crate) async fn run_tun_inbound(
handler: TunHandler,
prepared: PreparedDevice,
restart: mpsc::UnboundedReceiver<(TunHandler, PreparedDevice)>,
plane: PlaneCell,
flows: Tracker,
stop: CancellationToken,
);

The listener keeps the device’s own descriptor, and each runtime is handed a duplicate prepared before the change that started it was committed. run_tun_inbound loops:

  1. Build the connector factory: for each flow, AppConnector::new(plane, FlowContext { inbound_tag: tag, source: Some(ip) }, flows, None), where ip is the flow’s source.
  2. Run handler.inbound.run(prepared, make, token) under token = stop.child_token(), raced against restart.recv().
  3. If a restart arrives, cancel token and wait for run to return. If run returns by itself, because stop fired or the device is gone, there is no next runtime.
  4. handler.inbound.shutdown().await: wait, polling every RELEASE_POLL (20 ms) for up to RELEASE_TIMEOUT (3 s), until the old runtime’s device descriptor has closed and no TCP stream is left. When the wait runs out, it logs warn tun: device fd or <n> tcp flows still open after 3s and goes on.
  5. With a restart and stop not cancelled, take the new handler and prepared duplicate and loop. Otherwise return.

So a runtime is cancelled, and its descriptor and TCP flows are given up to RELEASE_TIMEOUT to be released, before the next one starts. The restart channel is unbounded, but only the actor writes to it, once per apply that swaps this listener. When the listener is dropped, its sender goes with it, recv returns None, and the loop ends the same way.

Swapping a TUN handler ends the device’s flows, so the plan marks the swap as a disruption, as it does a change of the device itself. etemenanki-app refuses a TUN change on reload; a restart applies it.

TunInbound::run serves one prepared device until its token is cancelled or the device is gone. Each runtime has its own max_flows semaphore, created by run, so a restart starts from a fresh count.

What the IP stack yields What run does
A TCP connection Takes a flow permit with try_acquire_owned, or drops the connection with debug tun: dropping a flow; the flow limit is reached. Serves it through serve_stream over a PassthroughCore (sniffing unless the handler turned it off); an error is logged at debug as tun: tcp flow <src> -> <dst> ended: <error>.
A UDP flow while udp is off Dropped.
A UDP flow from a source that has an association Handed to that association through its bounded channel (FLOW_QUEUE, 16); a flow that finds the channel full is dropped, and the client retransmits.
A UDP flow from a new source Takes a flow permit, or is dropped with the same debug line, and starts that source’s one association: a runtime over TunUdpCore. Its end is logged at trace as tun: udp association of <src> ended: <error>.
Any other protocol, ICMP included Dropped, trace tun: dropping a packet of an unsupported protocol.

The defaults a front end fills in come from protocols/src/tun/config.rs: DEFAULT_MTU 1500, DEFAULT_UDP_IDLE_TIMEOUT 60 s and DEFAULT_MAX_FLOWS 65,536. The stack, the associations and the TCP tracking are on TUN.

A TUN device admits no users: its user table is UserTable::None, every flow carries Principal::anonymous(), and the connector has no session. So TUN flows are not listed by Supervisor::sessions, Supervisor::close does not reach them, and they count toward no user’s usage. The flows are still registered and metered in the flow tracker, with no session. TCP flows through the device get their own silent-client watchdog in the protocols crate (serve_stream applies HANDSHAKE_TIMEOUT per step until its PassthroughCore is established and fails with tun: the client never spoke).

Listener::stop cancels the listener’s token.

Listener What stops What runs on
Stream The accept loop, which drops the listener: the port closes, a Unix socket file is removed Every accepted connection, under its session’s token
Hysteria 2 Admission: new clients are refused Every QUIC connection with a session, until it ends or its session is closed. The endpoint closes, and the port is released, after the last one.
TUN The runtime, whose token is a child of the listener’s: every flow on the device ends with it Nothing

Actor::commit stops a listener for Step::StopAccepting, when an inbound was removed or moved to another bind, and removes it from Actor::listeners with swap_remove. The listener is dropped right after, which also drops the TUN device’s own descriptor, the Hysteria 2 handle the actor held and the stream handler’s watch sender.

The plan decides whether the connections go too:

Change Steps Connections
Inbound removed StopAccepting, then CloseSessions Closed: Sessions::close(Scope::Inbound(tag)) cancels every session of the tag (removing_an_inbound_closes_its_sessions)
Inbound moved to another bind Bind, StopAccepting Not closed by the plan: sessions belong to the tag, which is unchanged
Inbound renamed on the same bind SwapHandler, CloseSessions of the old tag The old tag’s sessions are closed; the listener is kept

StopAccepting comes before CloseSessions, and Sessions::open checks the stop token under the same lock the sweep takes, so a connection being accepted as its inbound is removed either is swept or never opens a session.

Actor::commit reports the listener steps in its ApplyReport:

Field Filled from
swapped Every SwapHandler
restarted Every Disrupt
rebound A Bind whose tag the running spec already had on another bind; a bind for a new tag is not listed
removed Every CloseSessions

a_hysteria2_obfuscation_change_needs_allow_disruptive reads these: an allowed obfuscation change lists the inbound in restarted and not in rebound.

Supervisor::shutdown(grace) asks the actor, which:

  1. stops every listener, and cancels every balancer’s probe token;
  2. closes the task tracker and waits up to grace for every task on it to finish;
  3. cancels the root token. Every session token is its child, so every connection task still running is dropped, and every Hysteria 2 listener closes its endpoint at once;
  4. waits for the task tracker again, without a limit, and clears Actor::listeners;
  5. stops the sampler and the usage sink, which report once more now that every session has ended and folded in its last bytes.

etemenanki-app passes SHUTDOWN_GRACE, 5 s, when it receives SIGTERM or Ctrl-C. When the last Supervisor handle is dropped without a shutdown, the actor runs shutdown(Duration::ZERO) itself. TUN flows get no grace: step 1 stops their runtime.

Removal and drain grace timers run on the same task tracker and end when the root token is cancelled. So a pending timer keeps step 2 waiting for the whole grace, and step 3 ends it rather than waiting out its deadline (shutdown_ends_pending_grace_timers).

A protocol whose server is a ProxyCoreDecode core, served over TCP or a Unix socket, plugs in at these places:

  1. Spec. Add a variant to InboundProtocolSpec in supervisor/src/topology/spec_plan/inbound.rs. If it runs over a transport, add it to stream_transport in supervisor/src/build/inbound.rs.

  2. Validation. validate_inbound in supervisor/src/build/validate.rs already pairs every protocol other than Hysteria 2 and TUN with a TCP or Unix bind. Add the protocol’s own checks: a transport check, or a user set when it has no open mode, as Trojan, VLESS and VMess do. See Validation.

  3. Credentials. Return the credential kind that admits a user in credential_kind (supervisor/src/build/users.rs).

  4. User table. Add a UserTable variant, build it in user_table, and add the pair to UserTable::fits and UserTable::store_into.

  5. Handler. Add a StreamProtocol variant in supervisor/src/topology/inbound/mod.rs that holds the table in an ArcSwap (or in a type that swaps its own users, as VMess does), name it in StreamProtocol::name, and map the table to it in build_handler.

  6. Serving. Add the core to the established! list in supervisor/src/serve.rs and add a serve_connection arm that calls drive::<{ NewCore::<Principal>::BUF_SIZE }, _, _> with the core built from load_full() of the table, sniff and source, then the connector and the handshake permit. drive requires Target = Flow, Error = io::Error and TransportAddr = (), and the core’s is_established must turn true only once the request is parsed and accepted; the existing cores reach that with Timing::enter(Phase::Relay, …).

  7. Disruption. If some setting change cannot spare live connections, add it to swap_disruption in supervisor/src/topology/spec_plan/plan.rs (Planning and applying a change).

  8. Front ends. Lower the protocol from each front end that should offer it: etemenanki-app in app/src/lower.rs (etemenanki-app: from TOML to a spec) and katana in src/lower/inbound.rs (Lowering inbounds and outbounds).

A protocol that cannot be a core, as SOCKS cannot with its second UDP socket, needs its own driver in serve_connection, with its own handshake timeout, and a WireCounted stream so its session still counts its bytes.

Invariant Enforced by Pinned by
A refused apply leaves no listener behind, and serves nothing it bound. Binding and pairing happen in prepare; dropping a Pending closes its socket, removes its Unix file, closes its endpoint or device. Only Listener::start, in commit, serves. a_refused_spec_changes_nothing in supervisor/tests/hot_swap.rs
An unchanged bind keeps its socket across an apply. Listeners are found by BindSpec; a changed inbound gets prepare_swap and swap, never a new bind. a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections, a_hysteria2_obfuscation_change_needs_allow_disruptive (rebound is empty); a_reload_keeps_the_inbound_serving_on_its_udp_port in app/tests/integration/e2e_hysteria_inbound.rs
An apply that keeps an inbound and its users does not touch its connections. Step::Reuse: no swap; a user table is stored for a reused inbound only when its admissions changed (same_admissions). an_established_connection_survives_an_apply_that_keeps_its_inbound
A new socket is served by the handler current at its accept. handler.borrow().clone() per accept; swap only replaces the watch value. a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections
Stopping a stream listener stops accepting and nothing else. Session tokens are children of the root, not of the listener. Plan side: a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions in supervisor/tests/unit/plan.rs. No test moves a listener with a live connection.
No session opens on a stopped listener. Sessions::open checks stop under the registry lock; the accept loop’s select! is biased toward stop. a_stopped_listener_opens_no_session in supervisor/tests/unit/session.rs
Removing an inbound closes its sessions. Step::CloseSessions after StopAccepting. removing_an_inbound_closes_its_sessions
A Unix listener removes only the socket file it created. SocketFile compares device and inode before removing. Removal: socks_over_a_unix_socket_relays_and_cleans_up in app/tests/integration/e2e_unix.rs. The device and inode check has no dedicated test.
The accept loop never waits for a permit. try_acquire_owned; a refused socket is dropped. No dedicated test
A persistent accept error does not spin the loop, and a per-client one does not slow it. should_backoff_accept_error and ACCEPT_ERROR_BACKOFF accept_error_backoff_classification in supervisor/tests/unit/serve.rs
A client that never completes its request is disconnected. drive’s per-step HANDSHAKE_TIMEOUT until established, and the core’s Timing deadline The core side: timing_arms_handshake_once_then_idle_per_byte_event in protocols/tests/unit/core/mod.rs, and per-core tests such as handshake_deadline_fails_the_connection in protocols/tests/unit/trojan/core.rs. drive’s own timeout has no test.
A session’s wire counts what its client’s side carried. drive adds each step’s transport bytes; WireCounted and charged datagrams for SOCKS; Wire::quic for Hysteria 2 mux_sub_flows_and_their_session_count_exactly_what_moved, a_udp_association_counts_each_sub_link_and_charges_its_session, a_hysteria2_session_bills_its_quic_connections_bytes in supervisor/tests/tracking.rs
Routing sees the inbound tag and the client address captured at accept. FlowContext built by connector_for and Shared::connector inbound_tag_selects_the_route, source_cidr_matches_the_client_address in app/tests/integration/e2e_route_context.rs
A Hysteria 2 listener never serves two endpoints on one port. Changes are applied inside the running endpoint by its setters. a_reload_keeps_the_inbound_serving_on_its_udp_port; switching_the_obfuscation_closes_the_old_clients_and_admits_new_ones in protocols/tests/pipeline/hysteria.rs
Shutdown is bounded by its grace, Hysteria 2 handshakes included. The root cancel runs Hy2Inbound::shutdown under the draining run. shutdown_does_not_wait_out_a_stalled_hysteria2_handshake
A TUN runtime is cancelled, and its descriptor and TCP flows given up to 3 s to be released, before the next one starts. run_tun_inbound awaits run and then TunInbound::shutdown before it loops. No dedicated test
What happens Where Result Log
A TCP or UDP address is in use or not allowed listener::bind ApplyError::Bind; the apply changes nothing The front end’s report, for example etemenanki-app’s failed to start: inbound <tag>: binding <bind> failed: <OS error>
A Unix path holds something other than a socket listener::bind ApplyError::Bind … failed: <path> exists and is not a socket
A TUN device cannot be created, addressed or routed tun::open ApplyError::Bind The OS error, or route <net>/<prefix>: <error> for a route
A certificate or key that does not parse build_handler ApplyError::Build building inbound <tag> failed: <error>, for example hysteria2: the certificate file contains no certificates
A handler or table of the wrong kind Pending::pair, prepare_swap, prepare_users ApplyError::Build (unreachable after validation) building inbound <tag> failed: the handler is not of the kind its listener serves
A Salamander key shorter than MIN_PSK_LEN (4 bytes) check_obfs ApplyError::Build; validation refuses it first with salamander obfs key must be at least 4 bytes malformed: hysteria2 obfs psk is too short
The Hysteria 2 endpoint cannot be created on the socket Hy2Inbound::open (from_socket, quinn’s endpoint) ApplyError::Build building inbound <tag> failed: <error>
A TUN descriptor cannot be registered with the reactor, or the stack refuses its MTU TunInbound::prepare, from Pending::pair or prepare_swap ApplyError::Build; the MTU case is unreachable after validation building inbound <tag> failed: <error>
The listener’s TUN descriptor cannot be duplicated for a swap prepare_swap (try_clone) ApplyError::Build building inbound <tag> failed: <error>
accept fails with ConnectionAborted or Interrupted run_stream_inbound Accepts again at once debug accept error: <error>
accept fails otherwise run_stream_inbound Sleeps 100 ms, or ends if stopped meanwhile warn accept error, backing off 100ms: <error>
Live or handshake semaphore empty run_stream_inbound Socket dropped debug dropping inbound connection; live connection limit reached or … handshake limit reached
The listener stops as a socket is accepted Sessions::open returns None Loop ends, socket dropped None
Transport handshake fails or takes over 10 s serve_socket Socket and session dropped debug inbound transport failed: <error>
No runtime step for 10 s before establishment drive TimedOut debug <protocol> connection from <source> ended: inbound handshake timed out after 10s
A request not complete 10 s after its first byte The core Runtime error debug … ended: proxy core: client did not complete its request in time
SOCKS handshake not complete in 10 s SocksInbound::serve TimedOut debug socks connection from <source> ended: client did not complete its request in time
Any other runtime or driver error drive, SocksInbound::serve Connection dropped debug <protocol> connection from <source> ended: <error>
A session’s token is cancelled (see Tokens and tasks) spawn_until The task’s future is dropped None
A Hysteria 2 QUIC handshake fails Hy2Inbound::run Connection dropped debug hysteria2: a handshake failed: <error>
A Hysteria 2 client arrives at the limit or while draining Hy2Inbound::run QUIC refusal debug hysteria2: refusing a connection; the listener is full or stopping
The Hysteria 2 port is still held 3 s after the endpoint closed Hy2Inbound::shutdown Returns anyway warn hysteria2: <address> did not come free within 3s
A TUN runtime’s descriptor or TCP flows outlive 3 s TunInbound::shutdown The loop goes on warn tun: device fd or <n> tcp flows still open after 3s
A TUN swap finds the runtime task ended Listener::swap A new runtime task is spawned on the prepared duplicate None
A Unix socket file cannot be removed SocketFile::drop Ignored; the next bind takes over the stale socket debug could not remove <path>: <error>

The serve.rs lines are logged under the target etemenanki_supervisor::serve, the listener lines under etemenanki_supervisor::system::listener. Nothing in serve.rs logs a single connection above debug.

For a stream connection, cancellation is by drop. When a session’s token fires, spawn_until’s select! drops the connection’s future at whatever .await it is parked on: the client socket closes, the runtime and every outbound it opened are dropped, the connection’s permit handles are dropped, and the last Arc<Session> handle deregisters the session and folds its last bytes into its user’s account. No protocol-level goodbye is sent. Protocol and transport code must therefore keep its cleanup in Drop, never after an .await. A Hysteria 2 session is closed by the protocols crate instead, with DISCONNECT_CODE (see One session per QUIC connection).

Constant Value Defined in Scope
MAX_LIVE_CONNECTIONS_PER_INBOUND 4,194,304 supervisor/src/serve.rs Per stream listener, for its whole life
MAX_HANDSHAKES_PER_INBOUND 204,800 supervisor/src/serve.rs Per stream listener, for its whole life
ACCEPT_ERROR_BACKOFF 100 ms supervisor/src/serve.rs Per accept error that backs off
HANDSHAKE_TIMEOUT 10 s protocols/src/core/mod.rs drive‘s per-step watchdog, the cores’ Timing, the SOCKS handshake, TUN TCP flows
TRANSPORT_HANDSHAKE_TIMEOUT 10 s protocols/src/transports/accept.rs An inbound’s TLS handshake, WebSocket upgrade and HTTP/2 handshake (InboundTransport::accept)
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs The sniffing window, inside the handshake phase
Core BUF_SIZE 16,384 to 65,569 bytes each core Three buffers of this size per stream connection
DEFAULT_MAX_CONNECTIONS 262,144 protocols/src/hysteria/server/config.rs Front ends’ default for a Hysteria 2 inbound’s max_connections
DEFAULT_MAX_CIRCUITS 4,194,304 protocols/src/hysteria/server/config.rs Front ends’ default for max_circuits
CLOSE_CODE, DISCONNECT_CODE 0x100, 0x107 protocols/src/hysteria/server/inbound.rs QUIC application codes: endpoint closed, one connection closed
MIN_PSK_LEN 4 bytes protocols/src/hysteria/obfs.rs Shortest Salamander key check_obfs and Hy2Inbound::open accept
RELEASE_TIMEOUT, RELEASE_POLL 3 s, 20 ms protocols/src/tun/inbound.rs Each TUN runtime’s shutdown
DEFAULT_MTU, DEFAULT_UDP_IDLE_TIMEOUT 1500, 60 s protocols/src/tun/config.rs Front ends’ defaults for a TUN inbound’s mtu and udp_idle_timeout
DEFAULT_MAX_FLOWS 65,536 protocols/src/tun/config.rs Front ends’ default for a TUN inbound’s max_flows, per runtime
FLOW_QUEUE 16 protocols/src/tun/udp.rs New UDP flows queued toward one TUN association
MIN_TUN_MTU 1280 supervisor/src/build/validate.rs Smallest TUN MTU a spec may name; the IP stack’s own floor
SHUTDOWN_GRACE 5 s app/src/main.rs etemenanki-app’s grace for Supervisor::shutdown

None of the serve.rs constants is configurable. A Hysteria 2 inbound’s max_connections and max_circuits must be at least 1. The limits across the whole workspace are collected on Limits, timeouts and memory.

Test File What it pins
accept_error_backoff_classification supervisor/tests/unit/serve.rs Interrupted and ConnectionAborted do not back off; OutOfMemory and Other do.
a_stopped_listener_opens_no_session supervisor/tests/unit/session.rs Sessions::open with a cancelled stop token returns None and registers nothing.
cancelling_the_root_cancels_every_session supervisor/tests/unit/session.rs Every session token is a child of the root.
an_established_connection_survives_an_apply_that_keeps_its_inbound supervisor/tests/hot_swap.rs An apply that adds an outbound and a rule reports the inbound as reused, swaps and rebinds nothing, and a live SOCKS relay keeps echoing.
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections supervisor/tests/hot_swap.rs SOCKS to HTTP on one port is a swap with no rebind; an old relay keeps working, a new client speaks HTTP, and a socket accepted before the apply is answered as SOCKS.
a_refused_spec_changes_nothing supervisor/tests/hot_swap.rs A bind failure on one new inbound releases another bound in the same prepare; the plane, users and listeners are as before.
removing_an_inbound_closes_its_sessions supervisor/tests/hot_swap.rs The removed inbound’s connection closes and its port refuses; the other inbound’s connection is untouched.
a_hysteria2_obfuscation_change_needs_allow_disruptive supervisor/tests/hot_swap.rs Refused with ApplyError::Disruptive by default; with allow_disruptive the inbound is in restarted and not in rebound.
a_tun_device_change_needs_allow_disruptive supervisor/tests/hot_swap.rs A device change and a settings change are each refused with ApplyError::Disruptive without allow_disruptive. Needs CAP_NET_ADMIN; skipped when a device cannot be created.
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake supervisor/tests/hot_swap.rs shutdown(Duration::ZERO) finishes within 5 s while a client’s QUIC handshake is stuck.
shutdown_ends_pending_grace_timers supervisor/tests/hot_swap.rs Hour-long removal and drain timers do not hold shutdown up, and the session they spared is closed.
mux_sub_flows_and_their_session_count_exactly_what_moved supervisor/tests/tracking.rs drive feeds the session exactly the bytes the socket carried.
a_udp_association_counts_each_sub_link_and_charges_its_session supervisor/tests/tracking.rs WireCounted counts the SOCKS greeting, authentication and request, and charged datagrams add their payload.
a_hysteria2_session_bills_its_quic_connections_bytes supervisor/tests/tracking.rs A Hysteria 2 session counts more than its 20,000-byte payload each way: the QUIC connection’s UDP bytes.
a_protocol_change_on_the_same_bind_swaps_the_handler_and_keeps_the_listener, a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions, a_removed_inbound_stops_accepting_and_closes_its_sessions, an_inbound_renamed_on_the_same_bind_swaps_and_closes_the_old_tags_sessions supervisor/tests/unit/plan.rs Which steps each kind of inbound change plans for its listener and sessions.
a_tun_settings_change_on_the_same_device_swaps_and_disrupts, a_hysteria2_obfs_change_swaps_and_disrupts, a_hysteria2_masquerade_change_swaps_without_disrupting supervisor/tests/unit/plan.rs Which swaps are disruptions.
two_inbounds_may_not_share_a_bind, the_same_port_over_udp_is_another_bind, each_protocol_is_served_only_on_its_kind_of_bind, a_unix_listener_carries_only_the_plain_tcp_shape supervisor/tests/unit/validate.rs The validation that makes BindSpec a key and pairs protocols with binds.
a_connections_close_token_closes_that_connection_alone protocols/tests/pipeline/hysteria.rs Cancelling the token make_connector returned closes that connection and not its neighbour.
cancelling_run_stops_admitting_and_waits_for_the_live_connections protocols/tests/pipeline/hysteria.rs A stopped Hysteria 2 listener refuses newcomers promptly, keeps relaying for the live connection, and run returns once it ends.
a_new_connection_config_reaches_only_connections_that_arrive_after_it protocols/tests/pipeline/hysteria.rs A new connection config serves the connections that arrive after it (a new password, no UDP).
switching_the_obfuscation_closes_the_old_clients_and_admits_new_ones, a_too_short_obfuscation_key_is_refused_and_changes_nothing protocols/tests/pipeline/hysteria.rs set_obfs on a running endpoint, and its refusal of a short key.
socks_over_a_unix_socket_relays_and_cleans_up, http_connect_over_a_unix_socket_relays app/tests/integration/e2e_unix.rs A SOCKS round trip and an HTTP CONNECT over a Unix listener; the socket file is gone after SIGTERM.
a_reload_keeps_the_inbound_serving_on_its_udp_port app/tests/integration/e2e_hysteria_inbound.rs A reload that changes a Hysteria 2 inbound’s settings with a client connected leaves it serving on the same port. Skipped without the reference hysteria build.
a_routed_connect_is_answered_while_the_app_runs app/tests/integration/e2e_tun.rs The TUN inbound’s IP stack answers a routed TCP connect, and the interface and route are gone after the app exits. Linux only; skipped without CAP_NET_ADMIN.
inbound_tag_selects_the_route, source_cidr_matches_the_client_address app/tests/integration/e2e_route_context.rs The tag and the peer IP captured at accept reach the router.

Neither the admission semaphores nor drive’s silent-client timeout is exercised by a test at this revision. A change to either should come with one; a paused Tokio clock (tokio::time::pause) makes the 10 s watchdog testable without waiting. Run the supervisor’s tests with cargo test -p etemenanki-supervisor and the app’s with cargo test -p etemenanki-app; see Testing.