Skip to content

Listeners and the serve loop

Source files: 16 · checked against katana v3.0.1 · Etemenanki 596916d
  • katana/src/manager/transport.rs
  • katana/src/manager/proxy.rs
  • katana/src/manager/node.rs
  • katana/src/serve.rs
  • katana/src/inbound.rs
  • katana/src/connector.rs
  • katana/src/meter.rs
  • katana/tests/unit/serve.rs
  • katana/tests/unit/connector.rs
  • katana/tests/unit/e2e.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/config.rs
  • Etemenanki/protocols/src/hysteria/server/datagrams.rs
  • Etemenanki/concepts/src/runtime.rs

This page follows a node’s traffic from the bound socket to the end of each connection. It covers one listener generation: the TransportManager that owns a bound listener, the Scope every task of that generation runs under, the accept loop and its three admission semaphores, the ProxyManager that holds the node’s user table, and drive, the loop that runs one connection’s runtime and decides when it ends.

It is written for contributors who change src/serve.rs, src/manager/transport.rs or src/manager/proxy.rs, or who add a protocol to a stream node. What happens inside a flow once it reaches the connector (admission, routing, audit, metering) is on the admission and metering pages; how the node decides to build or refresh a generation is on the node manager page.

The serving layer does four things, and leaves everything else to the kernel runtime and to the connector:

Concern Owner What it decides
Binding and teardown src/manager/transport.rs → TransportManager Validate, bind, start serving; stop, drain and give the port back.
Admission of sockets and streams src/serve.rs → accept_loop, serve_socket; src/manager/proxy.rs → ProxyManager::accept_stream How many sockets a node holds, how many are in a transport handshake, how many streams are in a protocol handshake.
One connection’s life src/serve.rs → serve_stream, drive Which protocol core to build, the handshake deadline, the progress watchdog, retirement.
The user table src/manager/proxy.rs → Tables, ProxyManager::refresh What each new connection authenticates against, swapped without touching the listener.

The runtime itself (buffers, effects, the deadline timer) belongs to etemenanki-concepts, and the protocol cores with their per-phase deadlines belong to etemenanki-protocols; see Server runtime. katana only wraps them: it chooses the core, supplies the connector, and watches the runtime from outside.

A node serves one of two listener shapes, and they share very little:

Stream node Hysteria 2 node
Socket tokio::net::TcpListener std::net::UdpSocket, handed to quinn
Transport InboundTransport: TCP, TLS, WebSocket or gRPC QUIC, inside Hy2Inbound
Who accepts katana’s accept_loop The listener’s own Hy2Inbound::run
Who runs the runtimes katana, one drive per stream The listener, one runtime per proxy stream and, when UDP is enabled, one per connection’s datagrams
User table Tables::Stream(ArcSwap<StreamProtocol>) Tables::Hysteria, whose authenticator is replaced in place
How a retired user’s connection ends drive sees the lease cancelled and returns The user’s outbounds refuse to move (Gate), and admission refuses new flows

Everything a generation owns hangs off one TransportManager. The node manager holds at most one of them at a time.

flowchart TB
  TM["TransportManager"]
  SC["Scope: CancellationToken + TaskTracker"]
  PM["ProxyManager (shared Arc)"]
  AM["handshake-failure sampler"]
  AL["accept_loop or run_hysteria"]
  SS["serve_socket, one per socket"]
  ST["serve_stream, one per stream"]
  TB["Tables"]
  DP["Dispatcher + Admission"]
  TM --> SC
  TM --> PM
  SC --> AM
  SC --> AL
  SC --> SS
  SC --> ST
  PM --> TB
  PM --> DP

Every task of the generation, including the sampler, is spawned into the same Scope, so cancelling that one scope stops the listener and every connection it accepted. The scope’s token is a fresh root token, not a child of any node-level token: the owner must call TransportManager::shutdown (or drop the manager) for it to fire.

src/manager/transport.rs → TransportManager is one bound listener generation:

pub struct TransportManager {
accept: Scope,
proxy: Arc<ProxyManager>,
}
impl TransportManager {
pub async fn start(
node: &NodeInfo,
cert: &CertConfig,
listen_ip: &str,
enable_vless: bool,
sniff: bool,
users: &[UserInfo],
traffic: Arc<NodeTraffic>,
rules: Arc<RuleManager>,
router: Arc<Router<Outbound>>,
node_tag: CompactString,
hysteria: &HysteriaConfig,
) -> io::Result<Self>;
pub fn proxy(&self) -> &Arc<ProxyManager>;
pub async fn shutdown(self);
}
impl Drop for TransportManager {
fn drop(&mut self);
}

proxy() is how the node manager reaches ProxyManager::refresh for a user-set or node speed-limit change that does not need a new listener. A transport or protocol change from the panel, a config-file edit to the listen address, the certificate settings, enable_vless, disable_sniffing, the [node.hysteria] table or the routing rules, and a new outbound pool instead replace the whole TransportManager, which drops every connection on the node. The panel cannot show a local edit, so the node manager forces that new generation itself.

A bound socket, before it is served. The enum is private to src/manager/transport.rs:

enum Listener {
Stream(TcpListener, InboundTransport),
Datagram(std::net::UdpSocket, Hy2Inbound<UserTag>),
}
fn bind_listener(listen: &str, port: u16) -> io::Result<TcpListener>;
fn bind_datagram(listen: &str, port: u16) -> io::Result<std::net::UdpSocket>;

bind_listener binds a std::net::TcpListener, sets it non-blocking and converts it with TcpListener::from_std. bind_datagram binds a plain std::net::UdpSocket and hands it to the Hysteria listener as it is, because binding inside the listener would race a second socket against this one.

src/serve.rs → Scope pairs a cancellation token with a TaskTracker, so a teardown can both stop and wait:

#[derive(Clone, Default)]
pub struct Scope {
pub token: CancellationToken,
tasks: TaskTracker,
}
impl Scope {
pub fn new() -> Self;
pub async fn shutdown(&self);
}
pub fn spawn_scoped<F>(scope: &Scope, fut: F)
where
F: Future + Send + 'static,
F::Output: Send;

spawn_scoped spawns a task on the tracker that races fut against token.cancelled() in a tokio::select!. When the token fires, the select completes and fut is dropped, whatever it was awaiting. Scope::shutdown cancels the token, closes the tracker and awaits TaskTracker::wait, so it returns only when every task spawned under the scope has finished. Clones share the token and the tracker, which is how accept_loop, serve_socket and accept_stream spawn into the same scope.

src/manager/proxy.rs → ProxyManager owns one listener’s user table and the admission behind it. It is always held as Arc<ProxyManager>, and all mutation is interior:

pub enum Tables {
Stream(ArcSwap<StreamProtocol>),
Hysteria {
server: Hy2Inbound<UserTag>,
cfg: HysteriaConfig,
},
}
pub struct ProxyManager {
tables: Tables,
sniff: bool,
dispatcher: Arc<Dispatcher>,
traffic: Arc<NodeTraffic>,
preauth: Arc<Semaphore>,
handshake_failures: AtomicU64,
node_tag: CompactString,
}
impl ProxyManager {
pub fn new(
tables: Tables,
sniff: bool,
dispatcher: Arc<Dispatcher>,
traffic: Arc<NodeTraffic>,
node_tag: CompactString,
) -> Arc<Self>;
pub fn dispatcher(&self) -> Arc<Dispatcher>;
pub fn sniff(&self) -> bool;
pub async fn release_listener(&self);
pub fn note_handshake_failure(&self);
pub fn spawn_auth_monitor(self: &Arc<Self>, scope: &Scope);
pub fn accept_stream(
self: &Arc<Self>,
stream: TransportStream,
source: IpAddr,
session: Arc<OwnedSemaphorePermit>,
scope: &Scope,
);
pub fn refresh(&self, node: &NodeInfo, users: &[UserInfo], enable_vless: bool);
pub fn retire_all(&self);
}

The two Tables variants are two different mechanisms, not one mechanism with two payloads:

  • Tables::Stream holds the stream node’s StreamProtocol (src/inbound.rs) behind an ArcSwap. Each connection reads it once, with load_full, and builds its own protocol core over that snapshot. A refresh stores a whole new table.
  • Tables::Hysteria holds a clone of the live Hy2Inbound. The listener is the UDP socket, so rebuilding it would rebind the port and drop every connected client on each panel sync. A refresh replaces only its authenticator, through Hy2Inbound::set_authenticator. The clone is also how release_listener reaches the live endpoint at teardown. cfg is kept so the refresh builds the authenticator the same way start did.

sniff is kept for the same reason as cfg: every core built after a refresh has to sniff the way the first ones did.

src/inbound.rs → StreamProtocol is the table a stream node’s cores authenticate against:

pub enum StreamProtocol {
Vmess(Arc<AccountValidator<UserTag>>),
Vless(Arc<vless::Validator<UserTag>>),
Trojan(Arc<trojan::Validator<UserTag>>),
ShadowsocksLegacy(Arc<Resolved<UserTag>>),
Shadowsocks2022 {
config: Arc<Ss2022ServerConfig<UserTag>>,
validator: Option<Arc<ss_2022::Validator<UserTag>>>,
},
}

Each variant’s payload is an Arc, so building a core from it is a reference-count increment, not a copy. How the table is built from a panel’s node and user list is on the inbound and outbound construction page.

src/serve.rs holds the functions that move a socket from accept to the end of its last stream:

pub async fn accept_loop(
tcp: TcpListener,
transport: InboundTransport,
proxy: Arc<ProxyManager>,
scope: Scope,
);
async fn serve_socket(
transport: Arc<InboundTransport>,
sock: TcpStream,
peer: IpAddr,
proxy: Arc<ProxyManager>,
scope: Scope,
session: Arc<OwnedSemaphorePermit>,
stage: OwnedSemaphorePermit,
);
pub(crate) async fn serve_stream(
proxy: Arc<ProxyManager>,
protocol: Arc<StreamProtocol>,
stream: TransportStream,
source: Option<IpAddr>,
preauth: OwnedSemaphorePermit,
);
async fn drive<const BUF: usize, Core, S, Conn>(
stream: S,
core: Core,
connector: Conn,
preauth: OwnedSemaphorePermit,
retired: watch::Receiver<Option<CancellationToken>>,
) -> Result<(), Ended>
where
S: AsyncRead + AsyncWrite + Unpin,
Core: ProxyCoreDecode<Target = Flow<UserTag>, Error = io::Error, TransportAddr = ()>
+ Established,
Conn: Connector<Flow<UserTag>>,
Conn::Datagram: DatagramLink<Addr = Destination>;
pub async fn run_hysteria(
inbound: Hy2Inbound<UserTag>,
socket: std::net::UdpSocket,
dispatcher: Arc<Dispatcher>,
token: CancellationToken,
);

drive reports how a connection ended through a private enum, and needs one capability from each core through a private trait:

enum Ended {
Handshake(String),
Relay(String),
}
trait Established {
fn is_established(&self) -> bool;
}

The established! macro implements Established for TrojanCore, VlessCore, VMessCore, ShadowsocksCore and Ss2022Core (each over UserTag) by forwarding to the core’s own is_established. That method is true once the core’s phase is Relay or Closing: the request is parsed and the flow is open or opening. A core in its sniffing window is not yet established.

The connector reaches back to the connection through src/connector.rs → LeaseSlot:

pub type LeaseSlot = watch::Sender<Option<CancellationToken>>;
impl KatanaConnector {
pub fn new(
disp: Arc<Dispatcher>,
source: Option<IpAddr>,
lease: Option<Arc<LeaseSlot>>,
) -> Self;
}

serve_stream creates a watch::channel(None), gives the sender to its KatanaConnector and keeps the receiver for drive. The first flow the connector admits publishes that user’s lease into the slot with send_if_modified; later flows leave it alone, because every flow on one connection belongs to one user. The lease itself is the per-user CancellationToken held by Admission; see Admission.

sequenceDiagram
  participant C as Client
  participant AL as accept_loop
  participant SS as serve_socket
  participant PM as ProxyManager
  participant D as serve_stream and drive
  participant K as KatanaConnector
  C->>AL: TCP connect
  AL->>AL: try_acquire session permit, refuse if full
  AL->>AL: acquire stage permit, wait if full
  AL->>SS: spawn_scoped
  SS->>SS: InboundTransport accept (TLS, WebSocket, HTTP/2)
  Note over SS: sink yields a stream, stage permit released
  SS->>PM: accept_stream
  PM->>PM: try_acquire pre-auth permit, load_full table
  PM->>D: spawn_scoped
  C->>D: protocol request
  D->>K: connect(flow): admit, publish the lease
  Note over D: core established, pre-auth permit released
  K-->>D: outbound dialled
  D-->>C: relay until the runtime ends, the lease is cancelled or the watchdog fires

For a gRPC transport the middle of this diagram repeats: one socket yields one stream per HTTP/2 stream, and each yielded stream goes through accept_stream on its own.

accept_loop wraps the transport in an Arc and creates two semaphores that live exactly as long as the loop, so a new generation starts with fresh counts:

  • sessions, sized MAX_LIVE_CONNECTIONS_PER_NODE (65 536);
  • stages, sized MAX_TRANSPORT_STAGES_PER_NODE (2048).

It then loops on tcp.accept() against scope.token.cancelled(). For each accepted socket:

  1. Take a session place, or refuse. sessions.try_acquire_owned() never waits. If the semaphore is empty, the loop records a handshake failure, logs dropping connection; live connection limit reached at debug level, and drops the socket. Refusing is deliberate: waiting here would stop the loop from accepting at all and let the kernel’s listen backlog absorb the overload, which turns a saturated node into a silent one. The permit is wrapped in an Arc, because every stream the socket yields holds a clone of it.

  2. Take a transport-stage place, waiting if needed. stages.acquire_owned() is awaited inside a select! against the scope token, so a teardown is never stuck behind it. While it waits, the loop accepts nothing new; the listen backlog holds further connections until some socket gives its stage place back.

  3. Spawn serve_socket under the same scope with both permits.

An accept error is classified by should_backoff_accept_error:

pub(crate) fn should_backoff_accept_error(e: &io::Error) -> bool;

ConnectionAborted and Interrupted concern one connection; the loop logs them at debug level (accept error: ...) and continues at once. Every other kind, for example running out of file descriptors, is treated as a condition the next accept would hit too. The loop logs accept error, backing off 100ms: ... at warn level and sleeps ACCEPT_ERROR_BACKOFF (100 ms) in backoff_or_cancelled, which returns early if the scope is cancelled. The loop never exits on an error; it ends only when the scope is cancelled, and ending drops the listener.

serve_socket runs InboundTransport::accept (from protocols/src/transports/accept.rs) over the socket:

pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);

The transport calls sink once for TCP, TLS and WebSocket and returns right after. For gRPC it calls sink once per HTTP/2 stream and returns only when the HTTP/2 connection ends. The transport applies its own 10 s deadline (TRANSPORT_HANDSHAKE_TIMEOUT in accept.rs) to the TLS handshake, the WebSocket upgrade and the HTTP/2 preface.

The stage permit sits in a parking_lot::Mutex<Option<OwnedSemaphorePermit>> and is released at whichever comes first:

  • the first stream: the sink closure takes the permit out of the mutex, sets the yielded flag, and calls ProxyManager::accept_stream;
  • katana’s own TRANSPORT_HANDSHAKE_TIMEOUT (10 s): a select! against tokio::time::sleep takes the permit and then keeps awaiting the same accept future.

The second path is there for gRPC. Its transport yields a stream only when the client opens one, which an idle client may never do, so the timer rather than the first stream bounds how long such a socket counts as being in its transport handshake. After the timer the connection continues to be served and keeps its session place.

If accept returns an error and no stream was ever yielded, the socket failed its transport handshake and serve_socket records a handshake failure. An error after the first stream is the connection ending, not a handshake failing, so it is only logged (inbound transport ended: ..., debug).

For each yielded stream, accept_stream:

  1. takes a pre-auth place with preauth.try_acquire_owned(). The semaphore has MAX_PREAUTH_STREAMS_PER_NODE (512) places. If none is free, the stream is dropped at once, a handshake failure is recorded, and node <tag>: dropping a stream; pre-auth limit reached is logged at debug level. Dropping rather than queueing is what keeps a scan from growing tasks and buffers without bound;
  2. matches Tables::Stream. A Hysteria node never receives a transport stream; the other arm is a debug_assert!(false, ...) and a silent return in release builds;
  3. loads the current table with ArcSwap::load_full. The stream keeps that snapshot for its whole life, so a refresh that happens later does not change which table this connection authenticated against. A later stream on the same gRPC socket loads the table that is current when it arrives;
  4. spawns serve_stream under the scope. The spawned future owns the session Arc clone, so the socket’s session place is held until its last stream ends.

serve_stream builds one connection’s pieces and dispatches on the table variant:

Table variant Core Buffer size (BUF_SIZE)
StreamProtocol::Vmess VMessCore::new(accounts, now_unix, sniff, source) 32 KiB
StreamProtocol::Vless VlessCore::new(validator, sniff, source) 16 KiB
StreamProtocol::Trojan TrojanCore::new(validator, sniff, source) 16 KiB
StreamProtocol::ShadowsocksLegacy ShadowsocksCore::new(resolved, sniff, source) 20 KiB
StreamProtocol::Shadowsocks2022 Ss2022Core::with_system_clock(config, validator, sniff, source) 32 KiB

BUF_SIZE is each core’s associated constant and becomes drive’s BUF parameter; the runtime allocates three buffers of that size per connection (transport read, transport staging, outbound scratch). The connector is KatanaConnector::new(proxy.dispatcher(), source, Some(Arc::new(lease))), where source is the peer address the listener saw.

When drive returns, serve_stream maps the result:

drive result Effect
Ok(()) Nothing logged.
Err(Ended::Handshake(e)) Records a handshake failure; logs inbound handshake failed: <e> at debug level.
Err(Ended::Relay(e)) Logs connection ended: <e> at debug level.

drive builds ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress(), which turns the runtime into a Stream of Traffic deltas (concepts/src/runtime.rs). Polling that stream is what moves the connection; dropping it cancels the connection and everything it opened.

stateDiagram-v2
  [*] --> Handshake
  Handshake --> Established: core is_established
  Handshake --> EndedHandshake: runtime error or EOF
  Handshake --> EndedHandshake: HANDSHAKE_TIMEOUT
  Established --> Established: delta moved bytes, watchdog reset
  Established --> Finished: runtime stream ends
  Established --> Finished: lease cancelled
  Established --> EndedRelay: runtime error
  Established --> EndedRelay: PROGRESS_WATCHDOG
  EndedHandshake --> [*]
  Finished --> [*]
  EndedRelay --> [*]

The handshake phase. A client that connects and never speaks produces no runtime event at all, and the core arms its own handshake deadline only on the first event. So drive watches the handshake as a whole from outside: it polls runtime.next() inside tokio::time::timeout(HANDSHAKE_TIMEOUT, ...) until runtime.core().is_established(). The deadline starts when the stream reaches drive, and it covers the request and any sniffing window. The three ways out of this phase all become Ended::Handshake:

Cause Message
The runtime yields an error The error’s text
The runtime stream ends closed before the handshake completed
The deadline passes timed out after 10s

Once the core is established, drive drops the pre-auth permit. Only the handshake counts against MAX_PREAUTH_STREAMS_PER_NODE; an established connection counts only against the session semaphore.

The relay phase. drive then selects over three futures:

  • runtime.next(). Some(Ok(delta)) resets the watchdog if moved(&delta), meaning any of transport_rx, transport_tx, outbound_rx or outbound_tx is non-zero. A zero delta (a connect, a close, a deadline) does not count as progress. Some(Err(e)) returns Ended::Relay. None means the core finished: the connection ends cleanly with Ok(()).
  • until_retired(retired). It waits until the slot holds a lease, then awaits lease.cancelled_owned(). It uses borrow_and_update before changed, so a lease published during the handshake is not missed. Until a flow is admitted, the connection has no user to lose and the future never resolves; if the sender is gone it parks on std::future::pending. A cancelled lease returns Ok(()): the user was retired, so the connection’s runtime, its client stream and every outbound are dropped together.
  • The progress watchdog, a tokio::time::sleep(PROGRESS_WATCHDOG) that is reset on every delta that moved bytes. If it fires, drive returns Ended::Relay("nothing moved for 360s").

The watchdog is a backstop, not the idle timeout. The core’s own RELAY_IDLE_TIMEOUT (300 s) ends a quiet flow first, and gracefully. The watchdog gives up a connection that stops moving without ending, for example one whose client stopped reading while bytes are staged for it, so that the socket, the outbound and the three buffers are released. PROGRESS_WATCHDOG is RELAY_IDLE_TIMEOUT.saturating_add(Duration::from_secs(60)), so it never pre-empts the core’s graceful close.

A Hysteria node has no accept loop in katana. run_hysteria builds a connector factory and hands the socket to the listener:

let make = move |ip| KatanaConnector::new(dispatcher.clone(), Some(ip), None);

The factory’s shape is fixed by Hy2Inbound::run (protocols/src/hysteria/server/inbound.rs):

impl<T: Send + Sync + 'static> Hy2Inbound<T> {
pub async fn run<C, F>(
&self,
socket: std::net::UdpSocket,
make_connector: F,
token: CancellationToken,
) -> io::Result<()>
where
F: Fn(IpAddr) -> C + Send + Sync + 'static,
C: Connector<Flow<T>> + Send + 'static,
C::Future: Send,
C::Stream: Send,
C::Datagram: DatagramLink<Addr = Destination> + Send;
}

run calls the factory with the client’s address for each runtime it starts: one per proxy stream and one per connection’s UDP. The factory passes no lease slot, because katana owns no task per Hysteria connection that could watch one. A retired user’s flows end instead because their metered outbounds refuse to move (Gate::poll_open returns ConnectionAborted with the user was retired once the lease is cancelled; see metering), and because admission refuses the flows they open next.

Hy2Inbound::run serves until the token passed to it is cancelled or the endpoint closes. Every connection task lives in a JoinSet owned by the run future, so when the scope drops that future, every live QUIC connection goes with it. The listener’s own admission is fixed by katana in build_hysteria (src/inbound.rs):

Setting Value Effect at the limit
ListenerConfig::max_connections HY2_MAX_CONNECTIONS = 4096 The incoming connection is refused (incoming.refuse()), so the client learns at once.
ServerConfig::circuit_permits Semaphore::new(MAX_LIVE_CONNECTIONS_PER_NODE) = 65 536 A proxy stream is reset with H3_REQUEST_REJECTED, or the packet that would open a new UDP session is dropped. The permit is held by the task that runs the stream’s relay, or by the UDP session entry in Hy2UdpCore, so it covers the circuit’s whole life.

run returns an error only when the socket cannot become a quinn endpoint. run_hysteria then logs hysteria listener failed: <e> at error level and returns. Nothing in this layer restarts it, and the node manager keeps the TransportManager in place, so the node rebuilds its listener only on a later change that forces a new generation, such as a transport change from the panel or one of the config-file edits listed under TransportManager.

The listener’s internals, including its HTTP/3 authentication and stream classification, are on the Hysteria 2 server page.

TransportManager::start runs every step that can fail before it binds, so a bad configuration never leaves a half-bound listener:

  1. Stage the user set: build_user_entries, then traffic.prepare(entries).
  2. build_transport(node, cert): rejects REALITY, PROXY protocol, cert.reject_unknown_sni, ACME modes and unsupported transports, and reads the certificate files.
  3. build_protocol(node, &valid_users, enable_vless, tag_for): builds the StreamProtocol table.
  4. bind_listener(listen_ip, node.port).
  5. Tables::Stream(ArcSwap::from_pointee(protocol)) and Listener::Stream(tcp, transport).

Only after the bind succeeds does start:

  1. commit the staged user set with traffic.commit(prepared). Nothing is serving yet, so there is nobody to retire. NodeTraffic outlives generations, so unchanged users keep their counters across a rebuild; see traffic accounting;
  2. build the Dispatcher (router, audit rules, node tag, and a fresh Admission) and the ProxyManager;
  3. create the Scope and spawn the handshake-failure sampler into it;
  4. spawn accept_loop or run_hysteria into the scope.

An error at any step up to and including the bind returns from start with nothing bound and nothing committed: traffic.prepare only stages the set and does not touch the registry. start does not retry; the node manager decides what happens next. For a node’s first generation it retries the whole bootstrap, panel reads included, after a wait that starts at 1 s and doubles up to 60 s or the poll period, whichever is shorter. After a failed rebuild the node has no listener until the next poll cycle builds one.

TransportManager::shutdown is the supported teardown:

  1. self.accept.shutdown().await: cancel the scope and wait until every task under it is gone. The accept loop, the sampler, every serve_socket and every serve_stream are dropped; for a Hysteria node the run future is dropped with every connection it owned. Dropping the accept loop closes the TCP listener.

  2. self.proxy.retire_all(): cancel every user’s lease in Admission and empty the lease map, so any metered outbound that is still alive refuses to move another byte.

  3. self.proxy.release_listener().await: give the socket back in bounded time. This does nothing for a stream node. For a Hysteria node it calls Hy2Inbound::shutdown, which closes the endpoint with application code 0x100, waits up to DRAIN_TIMEOUT (3 s) for it to go idle, drops it, then tries to bind the same local address every RELEASE_POLL (20 ms) for up to RELEASE_TIMEOUT (3 s). If the port is still busy, it logs hysteria2: <addr> did not come free within 3s at warn level and returns.

Step 3 exists because quinn gives its socket back only when its driver task is polled with no connections and no endpoint handle left, about a millisecond after the last handle is dropped. The node manager binds the same port as soon as shutdown returns; without the wait, it would lose that race every time and the node would stay dark until its next poll cycle. The port itself is what Hy2Inbound::shutdown waits on, because it is the condition that has to be true.

impl Drop for TransportManager cancels accept.token and calls proxy.retire_all(). Neither TaskTracker nor CancellationToken does anything when dropped, so without this, one missed shutdown call would leave the listener, its accept loop and every live connection running for the life of the process. Drop cannot await, so it cannot wait for the tasks to finish or for a QUIC port to come free; it is a backstop, not a substitute for shutdown.

shutdown takes self by value, so Drop also runs at the end of every shutdown. Both of its actions are idempotent: the token is already cancelled and the lease map is already empty.

A user-set or node speed-limit change that needs no new listener goes through ProxyManager::refresh, which never touches the listener or a live connection:

  1. Stage the new user set and build the replacement: a new StreamProtocol for a stream node, or a new Authenticator for a Hysteria node. If the build fails, it logs proxy refresh build failed, keeping current: <e> at error level and returns with the running state untouched.
  2. Commit the registry through dispatcher.admission.commit(prepared). This refuses new flows of departed users and cancels the leases of departed users and of credentials rebound to a different uid, which ends their connections through drive or through Gate. A user whose speed limit changed keeps their connections; flows opened from then on use the new rate.
  3. Publish the table: ArcSwap::store for a stream node, Hy2Inbound::set_authenticator for a Hysteria node.

The registry goes first on purpose. Until the table swap, a departed user can still authenticate against the old table, but admission already refuses their flows. The other order would let a newly added user authenticate against the new table and then have every flow refused by a registry that has not heard of them. The full reasoning is on the runtime reload and admission pages.

ProxyManager::spawn_auth_monitor spawns, under the generation’s scope, a task that ticks every second (the immediate first tick is consumed). Each tick swaps handshake_failures to zero; if the count exceeds HANDSHAKE_FAILURE_ALERT_PER_SEC (10), it logs at warn level:

node <tag>: <n> inbound handshake failures in the last 1s (possible handshake scan/DoS or misconfigured clients)

It is detection only and never blocks anything. The counter is fed by note_handshake_failure at four places on the stream path:

Where What it counts
accept_loop A socket refused at the live-connection limit
serve_socket A transport failure before the first stream
ProxyManager::accept_stream A stream dropped at the pre-auth limit
serve_stream Any Ended::Handshake: protocol error, early close or HANDSHAKE_TIMEOUT

A Hysteria listener handles its own handshakes and logs failed ones at debug level (hysteria2: a handshake failed: ...); they do not reach this counter.

Invariant Enforced by Pinned by
A bad node config never half-binds. TransportManager::start runs every builder before bind_listener or bind_datagram, and commits traffic only after the bind. Builder rejections: tests/unit/inbound.rs (for example reality_rejected, hysteria_requires_a_certificate). The ordering itself has no dedicated test.
An accepted socket holds its session place until its last stream ends. The session permit is an Arc<OwnedSemaphorePermit>; serve_socket and every task spawned by accept_stream hold a clone. No dedicated test.
Overload at the session cap is refused, never queued. try_acquire_owned in accept_loop. No dedicated test.
A socket holds its transport-stage place for at most 10 s. The stage permit is taken by the first sink call or by the TRANSPORT_HANDSHAKE_TIMEOUT branch in serve_socket. No dedicated test.
Each stream authenticates against the table current at its arrival. ArcSwap::load_full in accept_stream; the snapshot moves into serve_stream. Indirectly by unchanged_user_survives_user_refresh in tests/unit/e2e.rs, where a live connection keeps relaying across a table swap.
A silent client ends exactly at HANDSHAKE_TIMEOUT. tokio::time::timeout around the handshake loop in drive. a_silent_client_is_dropped_at_the_handshake_deadline in tests/unit/serve.rs.
A pre-auth place is held only until the core is established. drop(preauth) right after the handshake loop in drive. Indirectly by the same test; no test fills the pre-auth semaphore.
A connection that moves nothing ends after PROGRESS_WATCHDOG. The watchdog Sleep, reset only by a delta for which moved is true. a_connection_that_stops_moving_is_dropped in tests/unit/serve.rs.
Retiring a user ends that user’s stream connections. The lease published through LeaseSlot, observed by until_retired. a_retired_users_connection_ends in tests/unit/serve.rs; the_lease_reaches_the_connection_and_goes_with_the_user in tests/unit/connector.rs; unchanged_user_survives_user_refresh in tests/unit/e2e.rs.
Retiring a user stops their Hysteria flows without disturbing other clients. Admission refuses new flows; Gate refuses to move bytes on a cancelled lease; the listener is not rebuilt. a_retired_user_stops_while_the_rest_keep_their_connections in tests/unit/e2e.rs.
A Hysteria refresh never rebinds the socket. Tables::Hysteria plus Hy2Inbound::set_authenticator. repeated_user_refreshes_never_disturb_a_live_connection in tests/unit/e2e.rs.
A teardown returns only after every task spawned into the generation’s scope is gone. Scope::shutdown (cancel, close, TaskTracker::wait). No test waits on the tracker directly; route_change_drops_connections in tests/unit/e2e.rs shows that a rebuild drops the open connection.
A QUIC port is free, or a warning is logged, when shutdown returns. ProxyManager::release_listener → Hy2Inbound::shutdown. No dedicated test.
A dropped TransportManager never strands its tasks. impl Drop for TransportManager. No dedicated test.
Hysteria flows are billed like stream flows. The same KatanaConnector, built per runtime by the factory in run_hysteria. a_hysteria_node_relays_and_meters in tests/unit/e2e.rs.

What ends each unit of work, and what it releases:

Unit Ends when Releases
accept_loop The scope is cancelled. Accept errors never end it. The TCP listener, and both semaphores with it.
serve_socket The transport returns: after the one stream for TCP, TLS and WebSocket, or when the HTTP/2 connection ends for gRPC; or the scope is cancelled. The stage permit, if still held, and its clone of the session permit.
serve_stream drive returns, or the scope is cancelled. The runtime (client stream, outbounds, buffers), the pre-auth permit if still held, and its clone of the session permit.
The sampler The scope is cancelled. Nothing else.
run_hysteria The scope token is cancelled, the endpoint closes, or run fails. Every QUIC connection task (owned by run’s JoinSet). The endpoint handle kept in Hy2Inbound is released by release_listener.

Cancellation always arrives as a dropped future: spawn_scoped selects each task against the scope token, so no task needs its own cancellation checks. A task that is waiting for anything (a transport handshake, a dial, a stage permit, a watchdog) is dropped at that point, and its permits are released by their destructors.

Constant Defined in Value Applies to
MAX_LIVE_CONNECTIONS_PER_NODE src/serve.rs 65 536 Accepted sockets alive at once per stream listener generation; also the size of a Hysteria listener’s circuit_permits.
MAX_TRANSPORT_STAGES_PER_NODE src/serve.rs 2048 Sockets inside their transport handshake at once. Past it, the accept loop waits.
TRANSPORT_HANDSHAKE_TIMEOUT src/serve.rs 10 s How long a socket holds its transport-stage place.
TRANSPORT_HANDSHAKE_TIMEOUT protocols/src/transports/accept.rs 10 s The transport’s own deadline on TLS, the WebSocket upgrade and the HTTP/2 preface.
ACCEPT_ERROR_BACKOFF src/serve.rs 100 ms Pause after an accept error other than ConnectionAborted or Interrupted.
MAX_PREAUTH_STREAMS_PER_NODE src/manager/proxy.rs 512 Streams in their protocol handshake at once per generation. Past it, streams are dropped.
HANDSHAKE_TIMEOUT protocols/src/core/mod.rs 10 s From drive’s start until the core is established.
RELAY_IDLE_TIMEOUT protocols/src/core/mod.rs 300 s The core’s own idle limit on a relaying flow.
PROGRESS_WATCHDOG src/serve.rs 360 s (RELAY_IDLE_TIMEOUT + 60 s) An established connection that moves no bytes.
HANDSHAKE_FAILURE_ALERT_PER_SEC src/manager/proxy.rs 10 Failures per 1 s sample above which the sampler warns.
HY2_MAX_CONNECTIONS src/inbound.rs 4096 Concurrent QUIC connections per Hysteria listener.
DRAIN_TIMEOUT protocols/src/hysteria/server/inbound.rs 3 s Wait for a closed endpoint to go idle.
RELEASE_TIMEOUT / RELEASE_POLL protocols/src/hysteria/server/inbound.rs 3 s / 20 ms Wait for the UDP port to become bindable again.

All semaphores are created per generation: the session and stage semaphores in accept_loop, the pre-auth semaphore in ProxyManager::new, and both Hysteria limits in build_hysteria and Hy2Inbound::run. A rebuild starts every count at zero, and a user refresh keeps the current counts.

The unit tests for this module live in tests/unit/serve.rs and are compiled into src/serve.rs as its tests module (#[path = "../tests/unit/serve.rs"]), so they can reach the private drive, Ended and should_backoff_accept_error.

Test What it pins
accept_error_backoff_classification Interrupted and ConnectionAborted do not back off; OutOfMemory and Other do.
a_silent_client_is_dropped_at_the_handshake_deadline On a paused clock, a Trojan core that receives no bytes ends with Ended::Handshake after exactly HANDSHAKE_TIMEOUT.
a_connection_that_stops_moving_is_dropped A passthrough flow whose far end floods a client that never reads ends with Ended::Relay containing nothing moved, no earlier than PROGRESS_WATCHDOG.
a_retired_users_connection_ends A connection keeps running while its lease is live, and returns Ok(()) within a second of the lease being cancelled.

The helpers in that file are useful for new tests of drive:

  • impl Established for PassthroughCore<UserTag>, so the kernel’s PassthroughCore (a core that opens a fixed flow on the first client bytes) can be driven;
  • once(outbound), a connector that hands one DuplexStream outbound to the first flow and refuses the rest;
  • flow(), a TCP flow to 192.0.2.1:80 for a user tag u with uid 1;
  • preauth(), a permit from a one-place semaphore.

The end-to-end tests in tests/unit/e2e.rs run a whole node against a fake panel: unchanged_user_survives_user_refresh, route_change_drops_connections, a_hysteria_node_relays_and_meters, a_retired_user_stops_while_the_rest_keep_their_connections and repeated_user_refreshes_never_disturb_a_live_connection. No katana test drives the caps themselves (the three semaphores, HY2_MAX_CONNECTIONS), the stage-permit release or the QUIC port release; a change to any of them should come with one.

Run the module’s tests from the katana repository with:

Terminal window
cargo test --locked serve::tests

A new protocol for stream nodes touches this page’s code in three places, besides its builder in src/inbound.rs:

  1. Add a StreamProtocol variant that holds the table as an Arc.
  2. Add the core to the established! list in src/serve.rs, so drive can tell when its handshake is over. The core must satisfy ProxyCoreDecode<Target = Flow<UserTag>, Error = io::Error, TransportAddr = ()>.
  3. Add a serve_stream arm that calls drive::<{ NewCore::<UserTag>::BUF_SIZE }, _, _, _> with the same connector, pre-auth permit and lease receiver as the others.

Nothing else changes: admission, the watchdogs, retirement and teardown apply to every arm of serve_stream alike.