Skip to content

Serving inbounds

Source files: 24 · checked against Etemenanki 596916d
  • Etemenanki/app/src/serve.rs
  • Etemenanki/app/src/connector.rs
  • Etemenanki/app/src/flow.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/inbound/tun.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/router.rs
  • Etemenanki/app/tests/unit/serve.rs
  • Etemenanki/app/tests/integration/e2e_unix.rs
  • Etemenanki/app/tests/integration/e2e_route_context.rs
  • Etemenanki/app/tests/integration/e2e_udp_route.rs
  • Etemenanki/app/tests/integration/e2e_hysteria_inbound.rs
  • Etemenanki/app/tests/integration/e2e_tun.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/protocols/tests/unit/socks/server.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/tests/unit/trojan/core.rs

app/src/serve.rs is where a bound listener turns into running connections. It owns the accept loop, the two per-inbound admission semaphores, the hand-off from an inbound transport to one task per byte stream, and the small driver that watches a protocol handshake before letting the per-connection runtime run free. Hysteria 2 and TUN inbounds, which own their socket or device outright, enter through two thin wrappers that hand the work to the protocols crate and then wait for the handle to come free.

Read this page before you add a stream protocol to the app, change how connections are admitted, or touch anything that runs under a generation’s cancellation token. How the listeners get bound in the first place is on Config to running pipeline; how a generation is replaced is on Generations and hot reload.

serve.rs and its two neighbours do the following:

Concern Where Notes
Bind every task to the generation app/src/serve.rs → spawn_scoped Cancelling the token drops the task’s future and everything it owns.
Accept TCP or Unix sockets app/src/serve.rs → StreamListener, run_stream_inbound Captures the peer IP as routing context, and as the one client a SOCKS UDP ASSOCIATE hears.
Admission run_stream_inbound Two Semaphores per inbound, taken with try_acquire_owned.
Transport fan-out serve_socket InboundTransport::accept yields one stream for TCP, TLS and WebSocket, and one per HTTP/2 stream for gRPC.
Protocol dispatch serve_connection SOCKS through its own driver; the other six protocols through drive.
Handshake watchdog drive HANDSHAKE_TIMEOUT between runtime events until the core is established.
Dialling app/src/connector.rs → AppConnector Routes each flow and opens it on the chosen outbound.
Hysteria 2 and TUN run_hysteria_inbound, app/src/inbound/tun.rs → run_tun_inbound Serve, then wait for the port or interface to be released.

It deliberately does not parse any protocol (the cores in etemenanki-protocols do), choose an outbound (the router does, through AppConnector), or decide when a generation ends (app/src/instance.rs does).

app/src/instance.rs → spawn_generation binds each built inbound and pairs its InboundKind with the Listener that bind_inbound produced. Each pair gets one task started with plain tokio::spawn, not spawn_scoped: these top-level tasks watch the token themselves, because they have cleanup to do after it fires, and the instance keeps their JoinHandles in Generation::accept_handles so a reload or shutdown can await them.

InboundKind BindSpec Bound handle Entry point
Stream(Arc<StreamInbound>) Tcp { host, port } StreamListener::Tcp serve::run_stream_inbound
Stream(Arc<StreamInbound>) Unix(PathBuf) StreamListener::Unix serve::run_stream_inbound
Hysteria2(Hy2Inbound<()>) Udp { host, port } std::net::UdpSocket serve::run_hysteria_inbound
Tun(TunInbound<()>) Tun(DeviceSpec) OwnedFd inbound::tun::run_tun_inbound

build_inbound pairs each kind with its bind spec, so the fall-through arm in spawn_generation, which logs inbound {tag}: listener does not match its kind, is unreachable in practice.

app/src/serve.rs
pub fn spawn_scoped<F>(token: CancellationToken, fut: F) -> JoinHandle<()>
where
F: Future + Send + 'static,
F::Output: Send,

The spawned task runs tokio::select! over token.cancelled() and fut. When the token fires first, the select! returns and fut is dropped in place. Dropping it closes the client socket, the per-connection runtime, every outbound connection it opened and every semaphore permit it held. This is how a reload aborts all old connections: the instance cancels the generation token and every scoped task unwinds by Drop, with no per-task bookkeeping.

The JoinHandle is returned but no caller keeps it: per-connection tasks are detached, and the token is the only handle on them.

app/src/serve.rs
pub enum StreamListener {
Tcp(TcpListener),
Unix {
listener: tokio::net::UnixListener,
path: PathBuf,
},
}
enum AcceptedSocket {
Tcp(tokio::net::TcpStream),
Unix(tokio::net::UnixStream),
}
impl StreamListener {
async fn accept(&self) -> io::Result<(AcceptedSocket, Option<IpAddr>)>;
fn release(self);
}
  • accept returns the peer’s IP for TCP (peer.ip(), the port is dropped) and None for a Unix socket, which has no IP peer.
  • release consumes the listener. For Unix, it drops the listener and then removes the socket file at path; a failure to remove is logged at debug level as could not remove {path}: {e} and is otherwise ignored. For Tcp, dropping the listener is all there is to do.

bind_inbound in app/src/instance.rs creates both variants from a non-blocking std listener. For a Unix path it takes over a stale socket file left by a crashed run or by the previous generation, and refuses a path that exists and is not a socket.

The accept loop is protocol-agnostic. Everything protocol-specific lives in the StreamInbound that build_inbound produced:

app/src/inbound/mod.rs
pub struct StreamInbound {
pub protocol: StreamProtocol,
pub transport: InboundTransport,
pub sniff: bool,
}
pub enum StreamProtocol {
Socks(SocksInbound<()>),
Http(Arc<HttpServerConfig<()>>),
Trojan(Arc<trojan::Validator<()>>),
Vless(Arc<vless::Validator<()>>),
Vmess(Arc<AccountValidator<()>>),
Shadowsocks(Arc<Resolved<()>>),
Ss2022 {
config: Arc<Ss2022ServerConfig<()>>,
validator: Option<Arc<ss_2022::Validator<()>>>,
},
}

Each variant holds what a per-connection core is built from; apart from Socks, the state sits behind an Arc, so building a core per connection costs a reference-count increment. The unit type () is the app’s per-user payload: the standalone app attaches nothing to a user. StreamProtocol::name returns the lowercase protocol name used in log lines ("socks", "http", "trojan", "vless", "vmess", "shadowsocks", "shadowsocks-2022").

transport is ignored for Unix listeners: build_inbound refuses a [inbound.stream] block on a Unix socket and stores InboundTransport::Tcp there, and serve_socket never consults it for an AcceptedSocket::Unix.

app/src/flow.rs
pub type Flow = etemenanki_protocols::flow::Flow<()>;
#[derive(Clone, Debug)]
pub struct FlowContext {
pub inbound_tag: CompactString,
pub source: Option<IpAddr>,
}
app/src/connector.rs
#[derive(Clone)]
pub struct AppConnector {
pub router: Arc<Router>,
pub ctx: FlowContext,
}
impl Connector<Flow> for AppConnector {
type Stream = OutboundStream;
type Datagram = FanOutLink;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow) -> ConnectFuture;
}

ConnectFuture is a boxed Future resolving to io::Result<Outbound<OutboundStream, FanOutLink>>. connect does one of two things:

  • UDP (flow.destination.network == DialNetwork::Udp): it dials nothing and returns Outbound::Datagram(FanOutLink::new(router, ctx, flow)) at once. The fan-out link routes every packet on its own; see Outbounds, fan-out and balancing.
  • Anything else: it calls router.route(route_target(&flow, &ctx)) synchronously, then returns the chosen outbound’s connect_stream(flow) future mapped into Outbound::Stream.

FlowContext exists so that two route matchers have something to match. route_target in app/src/router.rs builds the route target from the flow’s destination, sniffed domain and network, and adds two things from the context: ctx.inbound_tag for RouteMatch::InboundTag, and flow.source.or(ctx.source) for RouteMatch::SourceCidr. A source the flow carries itself (a QUIC client, a TUN host) wins over the one the listener captured. See Route model for the matchers.

flowchart TB
  gen["Generation token (CancellationToken)"]
  loop["run_stream_inbound (tokio::spawn)"]
  hy["run_hysteria_inbound (tokio::spawn)"]
  tun["run_tun_inbound (tokio::spawn)"]
  sock["serve_socket (spawn_scoped), one per accepted socket"]
  conn["serve_connection (spawn_scoped), one per yielded stream"]
  unixconn["serve_connection, inline for a Unix socket"]
  hyconns["per-connection tasks in a JoinSet"]
  tunflows["per-flow tasks in a JoinSet"]
  gen --> loop
  gen --> hy
  gen --> tun
  loop --> sock
  sock --> conn
  sock --> unixconn
  hy --> hyconns
  tun --> tunflows

The three entry tasks observe the token and return, and the instance awaits them. Everything under run_stream_inbound is scoped to the same token and is dropped when it fires. Hysteria 2 and TUN keep their per-connection work in a JoinSet owned by the run future, so returning from run drops those tasks too.

For a TCP socket, serve_socket and serve_connection are separate tasks: serve_socket runs the transport and spawns one serve_connection per stream the transport yields. For a Unix socket there is no transport, and serve_socket awaits serve_connection inline.

The sequence below follows a single TCP client of a VLESS-over-TLS inbound from accept to the end of the relay.

sequenceDiagram
  participant C as Client
  participant L as run_stream_inbound
  participant S as serve_socket task
  participant T as InboundTransport
  participant K as serve_connection task
  participant R as ProxyServerRuntime
  participant A as AppConnector
  C->>L: TCP connect
  L->>L: try_acquire_owned session, then handshake
  L->>S: spawn_scoped with FlowContext and both permits
  S->>T: accept(tcp, sink)
  T->>C: TLS handshake, within TRANSPORT_HANDSHAKE_TIMEOUT
  T->>K: sink(stream) spawns serve_connection
  K->>R: drive builds the runtime over VlessCore
  C->>R: request header
  Note over R: VlessCore pushes Effect::Open and enters Phase::Relay
  Note over K,R: core established, drive drops its handshake permit
  R->>A: connect(flow)
  A->>A: route_target, then router.route
  A-->>R: Outbound::Stream
  C->>R: payload both ways until the runtime ends
  R-->>K: None, drive returns Ok
  Note over K: session permit released when the task ends
app/src/serve.rs
pub async fn run_stream_inbound(
inbound: Arc<StreamInbound>,
tag: CompactString,
listener: StreamListener,
router: Arc<Router>,
token: CancellationToken,
)

The loop is a tokio::select! between token.cancelled() and listener.accept(). On cancellation it breaks out and calls listener.release(), so the port is closed, and a Unix socket file removed, by the time the instance’s await on this task’s handle returns. That ordering is what lets the next generation bind the same address.

Each accepted socket’s peer IP becomes FlowContext { inbound_tag: tag.clone(), source: peer }. It is carried rather than discarded because RouteMatch::SourceCidr has nothing to match without it, and because the SOCKS driver holds each UDP association to it (see serve_connection). The context is cloned into every stream a transport yields, so all gRPC streams on one connection share the same source.

app/src/serve.rs
const MAX_HANDSHAKES_PER_INBOUND: usize = 2048;
const MAX_LIVE_CONNECTIONS_PER_INBOUND: usize = 65_536;

run_stream_inbound creates both semaphores when it starts, so the counts are per inbound and per generation:

Semaphore Variable Capacity Held for
Live connections sessions MAX_LIVE_CONNECTIONS_PER_INBOUND (65,536) The whole life of the accepted socket: the transport, and every stream it yields until that stream’s task ends.
Handshakes handshakes MAX_HANDSHAKES_PER_INBOUND (2,048) The handshake stage. drive drops its reference as soon as the core reports established.

Both are taken with try_acquire_owned, in that order: session first, then handshake. The loop never waits for a permit. When either semaphore is empty, the loop logs at debug level and continues, which drops the freshly accepted socket (the client sees the connection close) and, when the handshake permit was the one missing, returns the session permit it had just taken:

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

Refusing at once, rather than queueing, keeps the accept loop responsive and keeps a burst from piling up sockets that nobody is serving.

Each permit is wrapped in Arc<OwnedSemaphorePermit> because one accepted socket can become many streams. serve_socket hands a clone of both Arcs to every serve_connection it spawns, and a permit returns to its semaphore when the last clone is dropped. serve_connection binds its session clone to a local (let _session = session;) so that it lives exactly as long as the connection task.

The live-connection cap is a guardrail, not a quota. The code comment on the constant gives the reasoning: a busy inbound legitimately carries tens of thousands of mostly idle connections, so the cap sits well above normal traffic, and it still leaves room under the process’s file-descriptor limit for the outbound socket each session may open. Operator-facing limits are on Limits and timeouts; how they compare with the other caps in the app is on Limits, timeouts and memory.

app/src/serve.rs
const ACCEPT_ERROR_BACKOFF: Duration = Duration::from_millis(100);
fn should_backoff_accept_error(e: &io::Error) -> bool;
async fn backoff_or_cancelled(token: &CancellationToken) -> bool;

should_backoff_accept_error sorts accept errors into two classes:

io::ErrorKind Meaning Handling
ConnectionAborted, Interrupted One client went away between the SYN and accept, or a signal interrupted the call. The listener is fine. Logged at debug as accept error: {e}, then the loop accepts again at once.
Anything else Typically resource exhaustion, such as running out of file descriptors, which would fail again immediately. Logged at warn as accept error, backing off 100ms: {e}, 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 token.cancelled() and returns true when the token won, in which case the loop breaks and releases the listener without finishing the sleep. An accept error never ends the loop on its own; only the token does.

app/src/serve.rs
async fn serve_socket(
inbound: Arc<StreamInbound>,
socket: AcceptedSocket,
ctx: FlowContext,
router: Arc<Router>,
token: CancellationToken,
session: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
)
  • AcceptedSocket::Unix: calls serve_connection directly on the Unix stream, with local_ip = None; the context’s source is None too. No transport runs and no keepalive is set.

  • AcceptedSocket::Tcp: records local_ip from tcp.local_addr() (the SOCKS driver puts it in its replies and, unless udp_bind is set, binds the UDP ASSOCIATE relay on it), then calls:

    protocols/src/transports/accept.rs
    impl InboundTransport {
    pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
    where
    F: FnMut(Accepted);
    }

    The sink closure is a spawn_scoped of serve_connection on the same generation token, with clones of the inbound, context, router and both permits. accept enables TCP keepalive on the socket, runs the transport handshake, and calls sink for every byte stream it yields:

    InboundTransport Streams yielded accept returns
    Tcp One, the socket itself Immediately after the one sink call
    Tls(ServerConfig) One, after the TLS handshake After the one sink call
    Ws { route, tls } One, after the (optionally TLS) WebSocket upgrade After the one sink call
    Grpc { paths, tls } One per HTTP/2 stream whose path matches the service; other paths are reset with REFUSED_STREAM When the HTTP/2 connection ends, or its liveness check judges the peer gone

    The transport’s own handshake (TLS, the WebSocket upgrade, the HTTP/2 handshake) is bounded by TRANSPORT_HANDSHAKE_TIMEOUT (10 s); a timeout fails with io::ErrorKind::TimedOut and the text tls handshake timed out, websocket handshake timed out or grpc handshake timed out. A transport error, including that timeout, ends serve_socket with a debug log line inbound transport failed: {e}. The transports themselves are described on TCP and TLS transports and WebSocket and gRPC transports.

app/src/serve.rs
async fn serve_connection<S>(
inbound: Arc<StreamInbound>,
stream: S,
local_ip: Option<IpAddr>,
ctx: FlowContext,
router: Arc<Router>,
session: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
) where
S: AsyncRead + AsyncWrite + Unpin + Send + 'static,

serve_connection builds one AppConnector { router, ctx: ctx.clone() } for the stream and then matches on inbound.protocol:

StreamProtocol Driver Core BUF (const generic)
Socks SocksInbound::serve none n/a
Http drive HttpCore::new(config, sniff, source) HttpCore::<()>::BUF_SIZE = MAX_HEAD = 64 KiB
Trojan drive TrojanCore::new(validator, sniff, source) TrojanCore::<()>::BUF_SIZE = 16 KiB
Vless drive VlessCore::new(validator, sniff, source) VlessCore::<()>::BUF_SIZE = 16 KiB
Vmess drive VMessCore::new(validator, now_unix, sniff, source) VMessCore::<()>::BUF_SIZE = 32 KiB
Shadowsocks drive ShadowsocksCore::new(resolved, sniff, source) ShadowsocksCore::<()>::BUF_SIZE = 20 KiB
Ss2022 drive Ss2022Core::with_system_clock(config, validator, sniff, source) Ss2022Core::<()>::BUF_SIZE = 32 KiB

sniff is inbound.sniff and source is ctx.source. BUF_SIZE is an associated constant on each core, and serve_connection passes it through as the runtime’s const generic, so each protocol gets buffers sized for its largest frame. The <()> turbofish is needed because the constant lives on the generic core type. The runtime allocates three buffers of BUF bytes per connection (transport read, transport staging, outbound scratch), so these constants set the size of a connection’s runtime buffers.

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

protocols/src/socks/server.rs
impl<T: Send + Sync + 'static> SocksInbound<T> {
pub async fn serve<S, C>(
&self,
mut stream: S,
local_ip: Option<IpAddr>,
source: Option<IpAddr>,
mut connector: C,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
C: Connector<Flow<T>>,
C::Datagram: DatagramLink<Addr = Destination>;
}

serve runs the SOCKS handshake under its own HANDSHAKE_TIMEOUT, failing with client did not complete its request in time when it expires, and then relays a CONNECT or drives a UDP ASSOCIATE; see SOCKS.

source is ctx.source: the peer IP that run_stream_inbound captured at accept, or None over a Unix socket. Besides going into each flow for routing, it is the only IP address a UDP ASSOCIATE hears (RFC 1928 §7). SocksInbound::associate builds an ExpectedSender (protocols/src/socks/server.rs) from source, the address and port the request names (its DST.ADDR and DST.PORT), and the relay’s bind IP, before it binds the relay socket:

Control connection Request names Association hears
TCP from source source with a non-zero port That IP and port only, from the start
TCP from source Anything else: source with port 0, another IP, an unspecified address or a domain source only; the first datagram forwarded pins the port
Unix socket (source = None) A specified IP and a non-zero port That IP and port only
Unix socket (source = None) An unspecified IP, a domain, or port 0 Nothing: refused

Addresses are compared in canonical form, so an IPv4-mapped IPv6 address is the IPv4 one. A datagram from any other sender is dropped unread, and replies go only to the pinned client. A source the request names over TCP that is not source is set aside rather than trusted or refused, because clients behind NAT name their LAN address. An association the relay could never hear, because udp_bind (or the TCP local IP) is in the other address family, is refused too; a relay on :: hears both families. Both refusals send the SOCKS5 reply 0x02 and end the connection with PermissionDenied, which serve_connection logs like any other driver error:

socks connection from None ended: socks: UDP associate over a unix socket must name its source address and port
socks connection from Some(203.0.113.7) ended: socks: UDP associate from an address family the relay is not bound in

Whatever the driver returns, an Err is logged at debug level with the protocol name and the client address:

vless connection from Some(203.0.113.7) ended: inbound handshake timed out after 10s

An Ok is silent. serve.rs logs nothing about an individual connection above debug level.

app/src/serve.rs
pub trait Established {
fn is_established(&self) -> bool;
}
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 the runtime in progress mode:

concepts/src/runtime.rs
impl<const BUF_SIZE: usize, Core, T, Conn>
ProxyServerRuntime<BUF_SIZE, Core, StreamTransport<T>, Conn, ProxyRunsQuiet>
{
pub fn new(transport: T, core: Core, connector: Conn) -> Self;
}
impl<const BUF_SIZE: usize, Core, Trans, Conn>
ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>
{
pub fn showing_progress(
self,
) -> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>;
}

In the default ProxyRunsQuiet mode the runtime is a Future that resolves once, at the end of the connection. showing_progress turns it into a Stream whose items are Result<Traffic, RuntimeError<Core::Error>>, one per unit of work. drive needs those intermediate yields: between two of them it can call runtime.core().is_established() and stop the watchdog at the right moment. The runtime model itself is described on The server runtime.

stateDiagram-v2
  [*] --> Handshake
  Handshake --> Handshake: Some(Ok), core not established
  Handshake --> Relay: core established, handshake permit dropped
  Handshake --> TimedOut: no event within HANDSHAKE_TIMEOUT
  Handshake --> Failed: Some(Err)
  Handshake --> Ended: None
  Relay --> Relay: Some(Ok)
  Relay --> Failed: Some(Err)
  Relay --> Ended: None
  TimedOut --> [*]
  Failed --> [*]
  Ended --> [*]

The two phases in detail:

  1. Handshake. While !runtime.core().is_established(), each runtime.next() is wrapped in tokio::time::timeout(HANDSHAKE_TIMEOUT, …).
    • A step Some(Ok(_)) loops.
    • A timeout returns io::ErrorKind::TimedOut with the text inbound handshake timed out after 10s.
    • Some(Err(e)) returns io::Error::other(e.to_string()).
    • None means the runtime finished before establishing, for example because the client closed or the core refused the request and finished; drive returns Ok(()).
  2. Relay. Once the core reports established, drive drops its handshake-permit reference (handshake.take() on an Option it wrapped the Arc in) and polls the runtime until the stream ends, converting any RuntimeError into io::Error the same way. The Traffic deltas are discarded: the standalone app does no accounting. katana’s serving layer drives the same runtime shape and meters those deltas; see Listeners and the serve loop.

is_established is true in the core’s Phase::Relay and Phase::Closing (protocols/src/core/mod.rs → Timing::is_established). A core enters Phase::Relay as soon as its request is parsed: for a plain TCP flow in the same event in which it pushes Effect::Open, and for a UDP or mux carrier (Trojan, VLESS, VMess) before any sub-flow is opened. “Established” therefore means the request is accepted, not that an outbound is connected: the outbound connect may still be in progress when drive stops the watchdog. A core that is still collecting a sniffing prefix (Phase::Sniff, bounded by SNIFF_TIMEOUT, 300 ms) is not established yet, so the watchdog and the handshake permit both cover sniffing.

A per-connection runtime delivers no event to its core until the client sends something, and the core’s own deadline (Timing) is armed only by its first byte event. A client that connects and says nothing would therefore sit forever inside the runtime. drive closes that gap from the outside, and the two timers complement each other:

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

Both use the same constant, HANDSHAKE_TIMEOUT in protocols/src/core/mod.rs, which is 10 s.

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

Hysteria 2 and TUN do not accept sockets; they own a UDP socket or a network device for the life of the generation. Their entry points share one shape: pass the inbound a factory that builds an AppConnector from a client address, run the inbound until the token fires, then call shutdown and wait for the handle to come free. They are started with tokio::spawn, and the instance awaits them, so a reload does not bind the next generation until the port or interface is released, or until shutdown gives up after RELEASE_TIMEOUT and logs a warning.

app/src/serve.rs
pub async fn run_hysteria_inbound(
inbound: Hy2Inbound<()>,
tag: CompactString,
socket: std::net::UdpSocket,
router: Arc<Router>,
token: CancellationToken,
)

The connector factory is move |ip| AppConnector { router: router.clone(), ctx: FlowContext { inbound_tag: tag.clone(), source: Some(ip) } }, and Hy2Inbound::run calls it with each QUIC connection’s remote IP, so every client carries its own address as routing context. Hy2Inbound::run(socket, make, token) serves until the token is cancelled or the endpoint closes; an error from it is logged at error level as hysteria2 inbound failed: {e}. Hy2Inbound::shutdown then closes the QUIC endpoint, waits up to DRAIN_TIMEOUT (3 s) for it to go idle, and polls every RELEASE_POLL (20 ms), for up to RELEASE_TIMEOUT (3 s), until the UDP port can be bound again; if it cannot, it logs the warning hysteria2: {local} did not come free within 3s and returns. The extra wait exists because dropping a quinn endpoint does not release the port while connections are still live, and the next generation would otherwise fail to bind.

Inside run, connection admission uses the listener’s own max_connections semaphore and refuses a QUIC handshake outright when it is full. See Hysteria 2 server.

Invariant Enforced by Pinned by
Cancelling a generation’s token ends every connection task it started. spawn_scoped selects on token.cancelled() and drops the future; Hysteria 2 and TUN keep their work in a JoinSet owned by run. Only indirectly: a_reload_rebinds_the_udp_port in app/tests/integration/e2e_hysteria_inbound.rs (a live connection does not keep the old generation’s port). No test reloads a stream inbound with a live connection.
A stream listener is closed, and its Unix socket file removed, before its accept task finishes. run_stream_inbound calls StreamListener::release after the loop; the instance awaits the task’s handle. socks_over_a_unix_socket_relays_and_cleans_up in app/tests/integration/e2e_unix.rs
A Hysteria 2 port or TUN interface is free again before the next generation binds. Hy2Inbound::shutdown and TunInbound::shutdown are awaited inside the entry task; the instance awaits the entry task. Hysteria 2: a_reload_rebinds_the_udp_port in e2e_hysteria_inbound.rs. TUN: no reload test; a_routed_connect_is_answered_while_the_app_runs in app/tests/integration/e2e_tun.rs only checks that the interface is gone after the process exits.
At most MAX_LIVE_CONNECTIONS_PER_INBOUND accepted sockets are live on one inbound. sessions semaphore, taken with try_acquire_owned at accept; every task serving the socket holds a clone of the permit. No dedicated test
At most MAX_HANDSHAKES_PER_INBOUND sockets on one inbound hold a handshake permit at once. handshakes semaphore, taken at accept; drive drops its reference once the core is established. No dedicated test
The accept loop never blocks on admission. try_acquire_owned instead of acquire; 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 plus ACCEPT_ERROR_BACKOFF. accept_error_backoff_classification in app/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 armed on the first byte event. The core side only: timing_arms_handshake_once_then_idle_per_byte_event in protocols/tests/unit/core/mod.rs; per-core tests such as handshake_deadline_fails_the_connection in protocols/tests/unit/trojan/core.rs. drive’s silent-client timeout has no test.
Routing sees the inbound tag and the client address. FlowContext built at accept, carried into AppConnector; route_target adds both. inbound_tag_selects_the_route and source_cidr_matches_the_client_address in app/tests/integration/e2e_route_context.rs
A SOCKS UDP association hears only the client its control connection came from. ExpectedSender, built from the source that serve_connection passes to SocksInbound::serve. udp_association_ignores_another_ip, udp_association_ignores_another_port_once_pinned and udp_association_over_a_unix_socket_needs_its_exact_source in protocols/tests/pipeline/socks.rs; ExpectedSender’s unit tests in protocols/tests/unit/socks/server.rs
A UDP flow is routed per packet, not dialled once. AppConnector::connect returns Outbound::Datagram(FanOutLink) for DialNetwork::Udp. one_association_routes_each_peer_separately in app/tests/integration/e2e_udp_route.rs
What happens Where Result Log
accept fails with ConnectionAborted or Interrupted run_stream_inbound Loop continues at once debug accept error: {e}
accept fails with any other kind run_stream_inbound Sleep 100 ms, or break if cancelled meanwhile warn accept error, backing off 100ms: {e}
Live-connection semaphore empty run_stream_inbound Socket dropped debug dropping inbound connection; live connection limit reached
Handshake semaphore empty run_stream_inbound Socket dropped, session permit returned debug dropping inbound connection; handshake limit reached
Transport handshake fails or exceeds 10 s serve_socket Socket dropped debug inbound transport failed: {e}
No runtime event for 10 s before establishment drive TimedOut, connection dropped debug {protocol} connection from {source:?} ended: inbound handshake timed out after 10s
Core or transport error drive Connection dropped debug {protocol} connection from {source:?} ended: {e}
SOCKS UDP ASSOCIATE with no client the relay can hear SocksInbound::serve Reply 0x02, connection dropped debug socks connection from {source:?} ended: socks: UDP associate …
Generation token cancelled every scoped task Future dropped mid-poll: sockets closed, outbound connections closed, permits returned none
Hysteria 2 or TUN run returns an error entry task shutdown still runs error hysteria2 inbound failed: {e} or tun inbound failed: {e}
Port or device not released within RELEASE_TIMEOUT Hy2Inbound::shutdown, TunInbound::shutdown The entry task returns anyway; the next generation’s bind may fail warn hysteria2: {local} did not come free within 3s or tun: device fd or {n} tcp flows still open after 3s
Unix socket file cannot be removed on release StreamListener::release Ignored; the next bind takes over a stale socket file debug could not remove {path}: {e}

Cancellation is by drop, not by a graceful close: a reload or shutdown ends every in-flight connection on the stream inbounds without a protocol-level goodbye. This is a functional property operators notice, and Hot reload documents it. Because a connection can be dropped at any .await, protocol and transport code must not rely on cleanup that runs only after an .await returns; put cleanup in Drop.

Constant Value Defined in Scope
MAX_LIVE_CONNECTIONS_PER_INBOUND 65,536 app/src/serve.rs Per stream inbound, per generation
MAX_HANDSHAKES_PER_INBOUND 2,048 app/src/serve.rs Per stream inbound, per generation
ACCEPT_ERROR_BACKOFF 100 ms app/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, the TUN inbound’s per-step watchdog
TRANSPORT_HANDSHAKE_TIMEOUT 10 s protocols/src/transports/accept.rs TLS, WebSocket upgrade, HTTP/2 preface
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs The cores’ Timing in relay, re-armed on every byte event
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs The sniffing window, inside the handshake phase
Core BUF_SIZE 16 KiB to 64 KiB each core Three buffers of this size per connection; see the table under serve_connection
DRAIN_TIMEOUT, RELEASE_TIMEOUT, RELEASE_POLL 3 s, 3 s, 20 ms protocols/src/hysteria/server/inbound.rs Hysteria 2 shutdown
RELEASE_TIMEOUT, RELEASE_POLL 3 s, 20 ms protocols/src/tun/inbound.rs TUN shutdown

None of the serve.rs constants is configurable.

A new protocol whose server is a ProxyCoreDecode core plugs into this file in four places:

  1. Add a variant to StreamProtocol in app/src/inbound/mod.rs holding the shared, Arc-wrapped state its core is built from, give it a name in StreamProtocol::name, and build it in build_inbound.

  2. Add the core to the established! list in app/src/serve.rs, so it implements Established.

  3. Add a match arm in serve_connection that calls drive::<{ NewCore::<()>::BUF_SIZE }, _, _> with the core built from (state, sniff, source), the connector and the handshake permit.

  4. Make sure the core’s Target is Flow, its Error is io::Error and its TransportAddr is (), as drive’s bounds require, and that its is_established turns true only once the request is parsed and accepted (for the existing cores, Timing::is_established, which the core reaches with Timing::enter(Phase::Relay, …)).

A protocol that cannot be expressed as a core, such as SOCKS with its second UDP socket, needs its own driver with its own handshake timeout.

Test File What it pins
accept_error_backoff_classification app/tests/unit/serve.rs Interrupted and ConnectionAborted do not back off; OutOfMemory and Other do.
socks_over_a_unix_socket_relays_and_cleans_up app/tests/integration/e2e_unix.rs A SOCKS round trip over a Unix listener, and removal of the socket file after SIGTERM.
http_connect_over_a_unix_socket_relays app/tests/integration/e2e_unix.rs An HTTP CONNECT driven through drive over a Unix stream.
inbound_tag_selects_the_route app/tests/integration/e2e_route_context.rs FlowContext::inbound_tag reaches the router.
source_cidr_matches_the_client_address app/tests/integration/e2e_route_context.rs The peer IP captured at accept reaches RouteMatch::SourceCidr, for both a matching and a non-matching range.
network_separates_tcp_from_udp app/tests/integration/e2e_route_context.rs TCP and UDP flows on one inbound take different routes.
udp_association_ignores_another_ip protocols/tests/pipeline/socks.rs A datagram from an IP other than the control connection’s is not forwarded, whether it reaches the relay before or after the client’s. Linux only.
udp_association_over_a_unix_socket_needs_its_exact_source protocols/tests/pipeline/socks.rs With source = None, a UDP ASSOCIATE naming 0.0.0.0:0 is refused with 0x02 and the connection closed, and one naming its exact address and port hears only that socket.
one_association_routes_each_peer_separately app/tests/integration/e2e_udp_route.rs Datagrams of one SOCKS UDP association are routed per destination through the fan-out link.
a_reload_rebinds_the_udp_port app/tests/integration/e2e_hysteria_inbound.rs A reload with a live QUIC connection releases and rebinds the same UDP port.
a_routed_connect_is_answered_while_the_app_runs app/tests/integration/e2e_tun.rs The TUN inbound’s IP stack answers a TCP connect while the app runs, and stops answering once the process has exited on SIGTERM. Linux only; skipped when creating a device is refused with PermissionDenied.
timing_arms_handshake_once_then_idle_per_byte_event protocols/tests/unit/core/mod.rs The core’s handshake deadline is armed once, and the idle deadline is refreshed on every relay byte event.
handshake_deadline_fails_the_connection protocols/tests/unit/trojan/core.rs A core whose handshake deadline passes fails with TimedOut.

The admission semaphores and drive’s silent-client timeout are not exercised by an app 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 app’s tests with cargo test -p etemenanki-app; the Hysteria 2 interoperability tests skip themselves when the reference hysteria binary is unavailable, and the TUN test needs CAP_NET_ADMIN to do anything. See Testing.