Listeners and the serve loop
Source files: 57 · checked against Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/system/mod.rsEtemenanki/supervisor/src/system/listener.rsEtemenanki/supervisor/src/topology/inbound/mod.rsEtemenanki/supervisor/src/build/inbound.rsEtemenanki/supervisor/src/lib.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/usage.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/topology/spec_plan/inbound.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/build/users.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/concepts/src/runtime.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/sniff/mod.rsEtemenanki/protocols/src/transports/accept.rsEtemenanki/protocols/src/transports/tls/config.rsEtemenanki/protocols/src/error.rsEtemenanki/protocols/src/socks/server.rsEtemenanki/protocols/src/http/core.rsEtemenanki/protocols/src/http/protocol.rsEtemenanki/protocols/src/trojan/core.rsEtemenanki/protocols/src/vless/core.rsEtemenanki/protocols/src/vmess/core.rsEtemenanki/protocols/src/ss_legacy/core.rsEtemenanki/protocols/src/ss_2022/core.rsEtemenanki/protocols/src/ss_2022/crypto.rsEtemenanki/protocols/src/ss_aead.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/hysteria/server/config.rsEtemenanki/protocols/src/hysteria/server/endpoint.rsEtemenanki/protocols/src/hysteria/obfs.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/protocols/src/tun/device.rsEtemenanki/protocols/src/tun/config.rsEtemenanki/protocols/src/tun/udp.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/supervisor/tests/unit/serve.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/validate.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/protocols/tests/pipeline/hysteria.rsEtemenanki/protocols/tests/unit/core/mod.rsEtemenanki/protocols/tests/unit/trojan/core.rsEtemenanki/app/tests/integration/e2e_unix.rsEtemenanki/app/tests/integration/e2e_hysteria_inbound.rsEtemenanki/app/tests/integration/e2e_tun.rsEtemenanki/app/tests/integration/e2e_route_context.rskatana/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.
Responsibilities
Section titled “Responsibilities”| 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).
Key types
Section titled “Key types”BindSpec: what is bound, and the listener’s identity
Section titled “BindSpec: what is bound, and the listener’s identity”#[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).
Handlers
Section titled “Handlers”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.
StreamInboundis what the accept loop serves every socket with.tagis what routes match on and what each session is filed under.transportis the layer under the protocol; a Unix listener carries none and never reads it.sniffsays whether a core reads a domain out of a flow addressed by IP.StreamProtocolholds what each connection’s core is built from. Every user-dependent part sits behind anArcSwap, 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 oneArc<AccountValidator>and swaps its user set itself (AccountValidator::set_users), which keeps its replay state. The type parameter isPrincipal, 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,shadowsocksorshadowsocks-2022.Ss2022Usersis 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.Hy2Handleris what a Hysteria 2 listener applies inside its one endpoint, because two QUIC endpoints cannot share a port.connectionis the per-connection config, with the user table (authenticator) and the circuit budget (circuit_permits) inside it.quicis the TLS material new handshakes are answered with.TunHandleris the TUN inbound itself, built whenever the inbound is new or changed.
Building handlers and user tables
Section titled “Building handlers and user tables”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 certificateshy2_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 masqueradeERROR 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 |
What a bind produces
Section titled “What a bind produces”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::acceptreturns the peer’s IP for TCP (peer.ip(); the port is dropped) andNonefor a Unix socket, which has no IP peer.- Dropping a
StreamListenerreleases 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. SocketFileis the filesystem entry a Unix listener created, with the file’s device and inode. ItsDropremoves the path only whilesymlink_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 atdebugascould not remove <path>: <error>and otherwise ignored._fileis declared afterlistener, so the descriptor closes before the file is removed.Bound::Tuncarries the device’s descriptor, which the listener keeps, andfirst, a duplicate for the first runtime.
Pending pairs, listeners and swaps
Section titled “Pending pairs, listeners and swaps”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
Pendingis a bound handle paired with the handler it will serve, with everything serving needs already built. It is the only thingListener::startaccepts, so starting cannot fail. - A
Listeneris a bound handle and the task serving it.tagis the tag of the inbound it currently serves, kept current byswap.bindis the key the actor finds it by.stopis a child of the supervisor’s root token, so shutdown reaches every listener. ListenerKindholds what the actor needs to reach the running task: thewatchsender 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.SwapandUsersare a handler and a user table already checked against the listener they go to. They are made in prepare and consumed in commit.
Shared, Connection and WireCounted
Section titled “Shared, Connection and WireCounted”#[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.
Data flow
Section titled “Data flow”Binding in prepare, serving in commit
Section titled “Binding in prepare, serving in commit”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
Tokens and tasks
Section titled “Tokens and tasks”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.
One accepted socket
Section titled “One accepted socket”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
Binding: listener::bind
Section titled “Binding: listener::bind”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:
check_platform(spec): off Linux, a spec with routes is refused withtun routes are installed only on Linux; add them with the OS route tool. Validation calls the same function for aTunSource::Create, so--testand a start agree.- 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. - Set it non-blocking and read back the name and interface index the kernel settled on.
- 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.
Unix socket files
Section titled “Unix socket files”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
bindnever meets its own socket; - anything else, a symlink included, is someone else’s file:
bindfails withio::ErrorKind::AlreadyExistsand<path> exists and is not a socket; NotFoundcontinues; any other error is returned.
- 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
std::os::unix::net::UnixListener::bind(path)creates the socket file.symlink_metadata(path)again records the new file’s device and inode in aSocketFile. If that fails, the file is removed and the error returned.set_nonblocking(true)andtokio::net::UnixListener::from_std. A failure here drops theSocketFile, 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 socketWith 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.sockA 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.
Dry runs bind nothing
Section titled “Dry runs bind nothing”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.
Pairing and starting
Section titled “Pairing and starting”Pending::pair
Section titled “Pending::pair”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:
- stores the Salamander key of the inbound’s
obfs, or none, into the inbound’s sharedSalamanderSwitch, so the endpoint obfuscates from its first packet.openfails on a keySalamander::newrefuses (shorter thanMIN_PSK_LEN, 4 bytes), althoughpairhas already refused one withcheck_obfs; - calls
endpoint::from_socket(socket, None, switch, gate)withgate = RecvGate::default(), a closed gate.from_socketsets the socket non-blocking, wraps it for quinn’s Tokio runtime, and wraps it again inSalamanderSocket::switchableeven 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.runopens the gate withRecvGate::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.
Listener::start
Section titled “Listener::start”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.
Swapping handlers and user tables
Section titled “Swapping handlers and user tables”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:
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>>;}Handler swaps
Section titled “Handler swaps”| 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).
User tables
Section titled “User tables”| 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.
The stream accept loop
Section titled “The stream accept loop”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.
Accept errors
Section titled “Accept errors”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.
Admission
Section titled “Admission”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 reacheddropping inbound connection; handshake limit reachedRefusing 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.
A session from the first moment
Section titled “A session from the first moment”With both permits taken, the loop:
- reads the current handler,
handler.borrow().clone(), which holds thewatchchannel’s read lock only while theArcis cloned; - opens the socket’s session,
shared.sessions.open(inbound.tag, peer, Wire::counted(), &stop). The session is registered, and listed bySupervisor::sessions, from this moment, with no user until its first flow binds one. Its transport handshake is closed with its inbound too.openchecks the stop token under the registry’s lock and returnsNoneonce 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); - builds a
Connectionwith the handler, the peer, the session ascarrier, clones ofstopandShared, and both permits; - spawns
connection.serve_socket(socket)withshared.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).
Connection::serve_socket
Section titled “Connection::serve_socket”serve_socket runs the transport over the accepted socket and serves every byte stream it yields.
AcceptedSocket::Unix:serve_connectionruns inline on the Unix stream, withlocal_ip = Noneand a connector for the carrier session withsource = None. No transport runs and no keepalive is set.AcceptedSocket::Tcp:local_ipis taken fromtcp.local_addr(). The SOCKS driver puts it in its replies and, unlessudp_bindis set, binds theUDP ASSOCIATErelay on it. ThenInboundTransport::accept(tcp, sink)runs.
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), andserve_connectionis spawned under the carrier’s token. - gRPC. Each HTTP/2 stream is a session of its own, opened through
Shared::connectorwith aWireof its own.
serve_connection: one driver per protocol
Section titled “serve_connection: one driver per protocol”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 10ssocks connection from None ended: client did not complete its request in timeThe SOCKS path
Section titled “The SOCKS path”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:
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 withclient did not complete its request in time. It then relays aCONNECTor drives aUDP ASSOCIATE; see SOCKS. - No runtime reports what the stream moved, so
serve_connectionwraps it inWireCountedwith 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).
drive: the handshake watchdog
Section titled “drive: the handshake watchdog”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 --> [*]
- Handshake. While the core is not established, each
runtime.next()is wrapped intokio::time::timeout(HANDSHAKE_TIMEOUT, …).Some(Ok(traffic))is counted and the loop goes on.- A timeout returns
io::ErrorKind::TimedOutwithinbound handshake timed out after 10s(the text is built fromHANDSHAKE_TIMEOUT.as_secs()). Some(Err(e))returnsio::Error::other(e.to_string()), so the log shows the runtime’s text.RuntimeErrorprefixes a transport failure withtransport:and a core error withproxy core:(for exampleproxy 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 versaandbytes staged toward a datagram transport without a peer. When each is raised is on The server runtime.Nonemeans the runtime finished before establishing, for example because the client closed or the core refused the request and finished.drivereturnsOk(()).
- Relay. Once the core is established,
drivedrops its handshake permit (handshake.take()on theOptionit holds it in) and polls the runtime until it ends, counting every step and converting aRuntimeErrorthe 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.
What established means
Section titled “What established means”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.
Two timers, one constant
Section titled “Two timers, one constant”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.
The Established trait
Section titled “The Established trait”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.
Hysteria 2 listeners
Section titled “Hysteria 2 listeners”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.
pub(crate) async fn run_hysteria_inbound( inbound: Hy2Inbound<Principal>, tag: Arc<ArcSwap<CompactString>>, endpoint: OpenedEndpoint, shared: Shared, stop: CancellationToken,);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);}From prepare to serving
Section titled “From prepare to serving”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:
- the token, once: cancelling it switches
runto draining; - a finished connection task, so finished tasks are reaped before new clients pile on;
- the next incoming client.
endpoint.accept()yieldingNonemeans 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.
One session per QUIC connection
Section titled “One session per QUIC connection”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:
- 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;
- 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; - 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.
Stopping and shutdown
Section titled “Stopping and shutdown”run_hysteria_inbound runs inbound.run(endpoint, make, stop) and races it against the session registry’s root token:
runreturns first. Cancellingstopmakesrundrain: 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,runreturns, andinbound.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 runsrunandinbound.shutdown()together:shutdowncloses the endpoint under the drainingrun, which ends every connection and handshake on it at once, andrunreturns with them.shutdown_does_not_wait_out_a_stalled_hysteria2_handshakepins 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.
TUN listeners
Section titled “TUN listeners”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.
pub(crate) async fn run_tun_inbound( handler: TunHandler, prepared: PreparedDevice, restart: mpsc::UnboundedReceiver<(TunHandler, PreparedDevice)>, plane: PlaneCell, flows: Tracker, stop: CancellationToken,);The restart loop
Section titled “The restart loop”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:
- Build the connector factory: for each flow,
AppConnector::new(plane, FlowContext { inbound_tag: tag, source: Some(ip) }, flows, None), whereipis the flow’s source. - Run
handler.inbound.run(prepared, make, token)undertoken = stop.child_token(), raced againstrestart.recv(). - If a restart arrives, cancel
tokenand wait forrunto return. Ifrunreturns by itself, becausestopfired or the device is gone, there is no next runtime. handler.inbound.shutdown().await: wait, polling everyRELEASE_POLL(20 ms) for up toRELEASE_TIMEOUT(3 s), until the old runtime’s device descriptor has closed and no TCP stream is left. When the wait runs out, it logswarntun: device fd or <n> tcp flows still open after 3sand goes on.- With a restart and
stopnot 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.
Inside one runtime
Section titled “Inside one runtime”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.
Flows without sessions
Section titled “Flows without sessions”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).
Stopping listeners and shutting down
Section titled “Stopping listeners and shutting down”What Listener::stop ends
Section titled “What Listener::stop ends”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
Section titled “Supervisor::shutdown”Supervisor::shutdown(grace) asks the actor, which:
- stops every listener, and cancels every balancer’s probe token;
- closes the task tracker and waits up to
gracefor every task on it to finish; - 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;
- waits for the task tracker again, without a limit, and clears
Actor::listeners; - 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).
Adding a stream protocol
Section titled “Adding a stream protocol”A protocol whose server is a ProxyCoreDecode core, served over TCP or a Unix socket, plugs in at these places:
-
Spec. Add a variant to
InboundProtocolSpecinsupervisor/src/topology/spec_plan/inbound.rs. If it runs over a transport, add it tostream_transportinsupervisor/src/build/inbound.rs. -
Validation.
validate_inboundinsupervisor/src/build/validate.rsalready 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. -
Credentials. Return the credential kind that admits a user in
credential_kind(supervisor/src/build/users.rs). -
User table. Add a
UserTablevariant, build it inuser_table, and add the pair toUserTable::fitsandUserTable::store_into. -
Handler. Add a
StreamProtocolvariant insupervisor/src/topology/inbound/mod.rsthat holds the table in anArcSwap(or in a type that swaps its own users, as VMess does), name it inStreamProtocol::name, and map the table to it inbuild_handler. -
Serving. Add the core to the
established!list insupervisor/src/serve.rsand add aserve_connectionarm that callsdrive::<{ NewCore::<Principal>::BUF_SIZE }, _, _>with the core built fromload_full()of the table,sniffandsource, then theconnectorand thehandshakepermit.driverequiresTarget = Flow,Error = io::ErrorandTransportAddr = (), and the core’sis_establishedmust turn true only once the request is parsed and accepted; the existing cores reach that withTiming::enter(Phase::Relay, …). -
Disruption. If some setting change cannot spare live connections, add it to
swap_disruptioninsupervisor/src/topology/spec_plan/plan.rs(Planning and applying a change). -
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 insrc/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.
Invariants
Section titled “Invariants”| 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 |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| 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).
Limits
Section titled “Limits”| 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.