Serving inbounds
Source files: 24 · checked against Etemenanki 596916d
Etemenanki/app/src/serve.rsEtemenanki/app/src/connector.rsEtemenanki/app/src/flow.rsEtemenanki/app/src/inbound/mod.rsEtemenanki/app/src/inbound/tun.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/router.rsEtemenanki/app/tests/unit/serve.rsEtemenanki/app/tests/integration/e2e_unix.rsEtemenanki/app/tests/integration/e2e_route_context.rsEtemenanki/app/tests/integration/e2e_udp_route.rsEtemenanki/app/tests/integration/e2e_hysteria_inbound.rsEtemenanki/app/tests/integration/e2e_tun.rsEtemenanki/concepts/src/runtime.rsEtemenanki/concepts/src/link.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/transports/accept.rsEtemenanki/protocols/src/socks/server.rsEtemenanki/protocols/tests/pipeline/socks.rsEtemenanki/protocols/tests/unit/socks/server.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/protocols/tests/unit/core/mod.rsEtemenanki/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.
Responsibilities
Section titled “Responsibilities”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).
Where serving starts
Section titled “Where serving starts”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.
Key types
Section titled “Key types”spawn_scoped
Section titled “spawn_scoped”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.
StreamListener and AcceptedSocket
Section titled “StreamListener and AcceptedSocket”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);}acceptreturns the peer’s IP for TCP (peer.ip(), the port is dropped) andNonefor a Unix socket, which has no IP peer.releaseconsumes the listener. ForUnix, it drops the listener and then removes the socket file atpath; a failure to remove is logged at debug level ascould not remove {path}: {e}and is otherwise ignored. ForTcp, 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.
StreamInbound and StreamProtocol
Section titled “StreamInbound and StreamProtocol”The accept loop is protocol-agnostic. Everything protocol-specific lives in the StreamInbound that build_inbound produced:
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.
FlowContext and AppConnector
Section titled “FlowContext and AppConnector”pub type Flow = etemenanki_protocols::flow::Flow<()>;
#[derive(Clone, Debug)]pub struct FlowContext { pub inbound_tag: CompactString, pub source: Option<IpAddr>,}#[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 returnsOutbound::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’sconnect_stream(flow)future mapped intoOutbound::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.
Data flow
Section titled “Data flow”The task tree
Section titled “The task tree”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.
One accepted socket
Section titled “One accepted socket”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
The accept loop
Section titled “The accept loop”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.
Peer address
Section titled “Peer 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.
Admission: two semaphores
Section titled “Admission: two semaphores”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 reacheddropping inbound connection; handshake limit reachedRefusing 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.
Accept errors and backoff
Section titled “Accept errors and backoff”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.
serve_socket: transport fan-out
Section titled “serve_socket: transport fan-out”async fn serve_socket( inbound: Arc<StreamInbound>, socket: AcceptedSocket, ctx: FlowContext, router: Arc<Router>, token: CancellationToken, session: Arc<OwnedSemaphorePermit>, handshake: Arc<OwnedSemaphorePermit>,)-
AcceptedSocket::Unix: callsserve_connectiondirectly on the Unix stream, withlocal_ip = None; the context’ssourceisNonetoo. No transport runs and no keepalive is set. -
AcceptedSocket::Tcp: recordslocal_ipfromtcp.local_addr()(the SOCKS driver puts it in its replies and, unlessudp_bindis set, binds theUDP ASSOCIATErelay on it), then calls:protocols/src/transports/accept.rs impl InboundTransport {pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>whereF: FnMut(Accepted);}The
sinkclosure is aspawn_scopedofserve_connectionon the same generation token, with clones of the inbound, context, router and both permits.acceptenables TCP keepalive on the socket, runs the transport handshake, and callssinkfor every byte stream it yields:InboundTransportStreams yielded acceptreturnsTcpOne, the socket itself Immediately after the one sinkcallTls(ServerConfig)One, after the TLS handshake After the one sinkcallWs { route, tls }One, after the (optionally TLS) WebSocket upgrade After the one sinkcallGrpc { paths, tls }One per HTTP/2 stream whose path matches the service; other paths are reset with REFUSED_STREAMWhen 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 withio::ErrorKind::TimedOutand the texttls handshake timed out,websocket handshake timed outorgrpc handshake timed out. A transport error, including that timeout, endsserve_socketwith a debug log lineinbound transport failed: {e}. The transports themselves are described on TCP and TLS transports and WebSocket and gRPC transports.
serve_connection: choosing a driver
Section titled “serve_connection: choosing a driver”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:
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 portsocks connection from Some(203.0.113.7) ended: socks: UDP associate from an address family the relay is not bound inWhatever 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 10sAn Ok is silent. serve.rs logs nothing about an individual connection above debug level.
drive: the handshake watchdog
Section titled “drive: the handshake watchdog”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:
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:
- Handshake. While
!runtime.core().is_established(), eachruntime.next()is wrapped intokio::time::timeout(HANDSHAKE_TIMEOUT, …).- A step
Some(Ok(_))loops. - A timeout returns
io::ErrorKind::TimedOutwith the textinbound handshake timed out after 10s. Some(Err(e))returnsio::Error::other(e.to_string()).Nonemeans the runtime finished before establishing, for example because the client closed or the core refused the request and finished;drivereturnsOk(()).
- A step
- Relay. Once the core reports established,
drivedrops its handshake-permit reference (handshake.take()on anOptionit wrapped theArcin) and polls the runtime until the stream ends, converting anyRuntimeErrorintoio::Errorthe same way. TheTrafficdeltas 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.
Why the watchdog lives here
Section titled “Why the watchdog lives here”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.
The Established trait
Section titled “The Established trait”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.
Inbounds that own their handle
Section titled “Inbounds that own their handle”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.
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.
pub async fn run_tun_inbound( inbound: TunInbound<()>, tag: CompactString, fd: OwnedFd, router: Arc<Router>, token: CancellationToken,)The connector factory is the same as for Hysteria 2. TunInbound::run calls it once per TCP flow and once per UDP association, with the client’s source IP. TunInbound::run(fd, make, token) serves the device until the token is cancelled or the device goes away; an error (for example from setting up the device or the IP stack) is logged at error level as tun inbound failed: {e}. TunInbound::shutdown then waits, polling every RELEASE_POLL (20 ms) for up to RELEASE_TIMEOUT (3 s), until the device descriptor has closed and no TCP stream is left, because the descriptor closes only after the IP stack’s task is gone and the next generation must be able to recreate the interface under the same name. If the wait runs out it logs the warning tun: device fd or {n} tcp flows still open after 3s and returns.
TUN TCP flows get the same silent-client treatment as drive: the inbound’s own driver, serve_stream in protocols/src/tun/inbound.rs, applies HANDSHAKE_TIMEOUT per runtime step until its PassthroughCore is established, and fails with tun: the client never spoke. See TUN.
Invariants
Section titled “Invariants”| 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 |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| 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.
Limits
Section titled “Limits”| 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.
Adding a stream protocol
Section titled “Adding a stream protocol”A new protocol whose server is a ProxyCoreDecode core plugs into this file in four places:
-
Add a variant to
StreamProtocolinapp/src/inbound/mod.rsholding the shared,Arc-wrapped state its core is built from, give it a name inStreamProtocol::name, and build it inbuild_inbound. -
Add the core to the
established!list inapp/src/serve.rs, so it implementsEstablished. -
Add a match arm in
serve_connectionthat callsdrive::<{ NewCore::<()>::BUF_SIZE }, _, _>with the core built from(state, sniff, source), theconnectorand thehandshakepermit. -
Make sure the core’s
TargetisFlow, itsErrorisio::Errorand itsTransportAddris(), asdrive’s bounds require, and that itsis_establishedturns true only once the request is parsed and accepted (for the existing cores,Timing::is_established, which the core reaches withTiming::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.