Skip to content

Life of a connection

Source files: 48 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/app/src/serve.rs
  • Etemenanki/app/src/connector.rs
  • Etemenanki/app/src/flow.rs
  • Etemenanki/app/src/router.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/inbound/tun.rs
  • Etemenanki/app/src/transport.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • Etemenanki/app/src/outbound/freedom.rs
  • Etemenanki/app/src/outbound/udp_fanout.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/socks/handshake.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/transports/keepalive.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/connector.rs
  • Etemenanki/protocols/src/hysteria/connection.rs
  • Etemenanki/protocols/src/wireguard/connector.rs
  • Etemenanki/protocols/src/wireguard/device.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/relay.rs
  • Etemenanki/environment/src/routing.rs
  • Etemenanki/environment/src/dial/tcp.rs
  • katana/src/serve.rs
  • katana/src/connector.rs
  • katana/src/manager/proxy.rs
  • katana/src/meter.rs
  • katana/src/router.rs
  • katana/src/outbound/mod.rs

This page follows one TCP connection through etemenanki-app, from the moment the listener accepts the socket to the moment the last byte drains and the task ends. It names every function, type and effect the connection passes through, in the order it meets them, and links each step to the page that covers that component in depth.

Read it first if you are about to change the serving path, a protocol core, the connector or an outbound: it shows which task owns what, where each timeout sits, and which event the core sees at each point. The later sections cover what changes for UDP, for the two inbounds that do not accept TCP sockets (Hysteria 2 and TUN), and for katana, which runs the same kernel behind its own connector.

  1. Accept and admission. run_stream_inbound accepts the socket and takes two semaphore permits without waiting, or drops the socket.
  2. Transport accept. serve_socket runs InboundTransport::accept: TLS, the WebSocket upgrade or the HTTP/2 preface, bounded by TRANSPORT_HANDSHAKE_TIMEOUT. gRPC yields one byte stream per HTTP/2 stream.
  3. Core selection. serve_connection builds an AppConnector and picks the protocol’s core. SOCKS goes to its own driver instead.
  4. Drive and the handshake watchdog. drive wraps the core in a ProxyServerRuntime and applies HANDSHAKE_TIMEOUT until the core is established.
  5. Request parsing. The core parses and authenticates the request and pushes Effect::Open with a Flow.
  6. Routing. While applying Open, the runtime calls AppConnector::connect, which turns the Flow and its FlowContext into a RouteTarget and asks the Router for an outbound.
  7. Dialing. The outbound’s connect_stream future dials: directly (freedom), through a proxy client runtime over a TransportConnector, through WireGuard, or through Hysteria 2.
  8. Connected. The runtime delivers Event::Connected; the core stages its reply and the forwards queued behind Open go out.
  9. Relay. Client bytes become Forward effects; outbound reads become Event::Outbound and are staged back toward the client.
  10. Half-close and completion. Each EOF becomes a half-close on the other side; when both have ended the core pushes Finish, staging drains, and the task ends, dropping its clones of the permits.
sequenceDiagram
  participant C as Client
  participant L as Accept loop
  participant T as serve_socket
  participant D as drive
  participant R as ProxyServerRuntime
  participant K as Core
  participant A as AppConnector
  participant X as Outbound
  C->>L: TCP connect
  L->>L: try_acquire session and handshake permits
  L->>T: spawn_scoped(serve_socket)
  T->>C: TLS, WebSocket upgrade or h2 preface
  T->>D: sink(TransportStream), spawn_scoped(serve_connection)
  D->>R: new(stream, core, connector).showing_progress()
  C->>R: request bytes
  R->>K: Event::Transport(unparsed)
  K-->>R: Effect::Open(key, Flow), Forward(rest)
  Note over D,K: is_established() is true, drive drops its handshake permit
  R->>A: connect(flow) while applying Open
  A->>A: route_target, Router::route
  A-->>R: boxed connect_stream future
  R->>X: poll the future: resolve, dial, upstream handshake
  X-->>R: Outbound::Stream
  R->>K: Event::Connected
  K-->>R: stage reply
  R->>X: queued Forward (poll_write, poll_flush)
  R->>C: staged reply
  loop Relay
    C->>R: bytes into the read buffer
    R->>K: Event::Transport
    K-->>R: Forward(range)
    R->>X: poll_write
    X->>R: bytes into scratch
    R->>K: Event::Outbound
    K-->>R: stage
    R->>C: poll_write from staging
  end
  C->>R: EOF
  R->>K: Event::TransportEof
  K-->>R: Shutdown(key)
  R->>X: poll_shutdown
  X->>R: EOF
  R->>K: Event::OutboundEof
  K-->>R: ShutdownTransport, Finish
  R->>C: drain staging, then shut down the write side
  R-->>D: stream ends
  D-->>L: task ends, permits drop

The whole connection lives in very few tasks. There is no task per direction and no channel between the client side and the outbound side: a ProxyServerRuntime owns the client transport, the core and every outbound the core opens, and moves bytes between them itself.

Task Spawned by Owns Ends when
run_stream_inbound app/src/instance.rs → spawn_generation, one per inbound the listener, both semaphores the generation’s CancellationToken is cancelled
serve_socket the accept loop, through spawn_scoped the accepted socket until the transport yields; for gRPC, the HTTP/2 connection the transport has yielded its stream, or the HTTP/2 connection ends
serve_connection serve_socket’s sink, through spawn_scoped (awaited inline for a Unix socket) the ProxyServerRuntime: transport, core, outbounds, three buffers the runtime completes or fails, or the generation is cancelled

spawn_scoped races the future against token.cancelled(), so cancelling a generation drops every future above, and with it every socket it owns. See Generations and hot reload.

A few outbounds rely on a background task: every dialed gRPC transport spawns the driver of its own HTTP/2 connection, the Hysteria 2 client runs an HTTP/3 driver (and, when the server relays UDP, a datagram pump) for the QUIC connection all its flows share, and the WireGuard outbound runs one device driver per tunnel. A proxy client runtime (ProxyClientRuntime) itself is polled inside the connection’s task, as the outbound’s AsyncRead and AsyncWrite.

The serving path in app/src/serve.rs:

pub fn spawn_scoped<F>(token: CancellationToken, fut: F) -> JoinHandle<()>
where
F: Future + Send + 'static,
F::Output: Send;
pub async fn run_stream_inbound(
inbound: Arc<StreamInbound>,
tag: CompactString,
listener: StreamListener,
router: Arc<Router>,
token: CancellationToken,
);
async fn serve_socket(
inbound: Arc<StreamInbound>,
socket: AcceptedSocket,
ctx: FlowContext,
router: Arc<Router>,
token: CancellationToken,
session: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
);
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;
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;
pub trait Established {
fn is_established(&self) -> bool;
}

What a decoded connection is, and where it entered (protocols/src/flow.rs, app/src/flow.rs):

pub struct Flow<T> {
pub destination: Destination,
pub user: NetworkUser<T>,
pub sniffed: Option<SniffedBehavior>,
pub source: Option<IpAddr>,
}
pub type Flow = etemenanki_protocols::flow::Flow<()>;
pub struct FlowContext {
pub inbound_tag: CompactString,
pub source: Option<IpAddr>,
}

The connector every stream inbound dials through (concepts/src/link.rs, app/src/connector.rs):

pub trait Connector<Target> {
type Stream: AsyncRead + AsyncWrite + Unpin;
type Datagram: DatagramLink;
type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;
}
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;
}

Routing and dialing (app/src/router.rs, environment/src/routing.rs, app/src/outbound/mod.rs):

pub fn route_target<'a>(flow: &'a Flow, ctx: &'a FlowContext) -> routing::RouteTarget<'a>;
pub fn route(&self, target: RouteTarget<'_>) -> Arc<Outbound>;
pub fn connect_stream(&self, flow: Flow) -> StreamFuture;
pub fn connect_datagram(&self, flow: Flow) -> DatagramFuture;

The runtime and the core contract it drives (concepts/src/runtime.rs, concepts/src/core.rs):

pub struct ProxyServerRuntime<const BUF_SIZE: usize, Core, Trans, Conn, Mode = ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Conn: Connector<Core::Target>;
pub fn new(transport: T, core: Core, connector: Conn) -> Self;
pub fn showing_progress(
self,
) -> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>;
fn handle(
&mut self,
event: Event<'_, Self>,
effects: &mut Effects<'_, Self>,
) -> Result<usize, Self::Error>;

run_stream_inbound loops on tokio::select! over the generation token and StreamListener::accept. A StreamListener is either a TcpListener or a Unix listener; only the TCP one reports a peer address, which becomes FlowContext.source so that source_cidr rules have something to match.

Admission happens before any byte is read, and it never waits:

  1. sessions.clone().try_acquire_owned() takes a place from the per-inbound live-connection semaphore (MAX_LIVE_CONNECTIONS_PER_INBOUND, 65,536).
  2. handshakes.clone().try_acquire_owned() takes a place from the per-inbound handshake semaphore (MAX_HANDSHAKES_PER_INBOUND, 2,048).

If either semaphore is full, the socket is dropped on the spot, together with any permit already taken, and the loop logs at debug level (dropping inbound connection; live connection limit reached or … handshake limit reached). Refusing instead of waiting keeps the loop accepting, so an overloaded inbound sheds connections rather than letting the kernel’s listen backlog fill silently.

The live-connection permit travels as an Arc<OwnedSemaphorePermit>. A socket can yield several byte streams (gRPC yields one per HTTP/2 stream), each served by its own task, and every one of those tasks keeps a clone of the live-connection permit: the place returns to the semaphore only when the last task serving the socket has ended.

An accept error is classified by should_backoff_accept_error: ConnectionAborted and Interrupted are retried at once, anything else (running out of descriptors, for example) is logged at warn level and followed by an ACCEPT_ERROR_BACKOFF (100 ms) sleep that is itself cancellable. When the token is cancelled the loop breaks and StreamListener::release removes a Unix socket file.

Deeper: Serving and Limits.

serve_socket hands a TCP socket to the inbound’s InboundTransport (protocols/src/transports/accept.rs):

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

accept first enables TCP keepalive on the socket (set_keepalive: first probe after TCP_KEEPALIVE_IDLE, 120 s, then every TCP_KEEPALIVE_INTERVAL, 30 s, giving up after TCP_KEEPALIVE_RETRIES, 3), then runs the transport’s own handshake inside within, which applies TRANSPORT_HANDSHAKE_TIMEOUT (10 s) and fails with tls handshake timed out (or websocket …, grpc …):

InboundTransport Handshake under the 10 s limit Streams yielded to sink
Tcp none exactly one, at once
Tls(ServerConfig) the TLS accept exactly one
Ws { route, tls } optional TLS, then the WebSocket upgrade (path and Host check) exactly one
Grpc { paths, tls } optional TLS, then the HTTP/2 server handshake one per accepted HTTP/2 stream on a known path, until the connection ends

For gRPC, serve_h2 keeps driving the HTTP/2 connection after the handshake. Each accepted stream whose path GrpcPaths::classify recognises (/{service}/Tun or /{service}/TunMulti) is answered with the gRPC response headers and passed to sink as a TransportStream::Grpc; any other path is reset with REFUSED_STREAM. Liveness watches the connection and gives up on a peer that has gone idle or stopped answering pings.

The sink in serve_socket spawns serve_connection for each yielded stream, under the same token, with the socket’s permits and its local_addr IP (which only the SOCKS driver uses). A Unix socket skips the transport entirely: the config builder refuses any stream network or security on a Unix listener (inbound <tag>: protocol http over a unix socket does not support stream network "ws"), so serve_socket awaits serve_connection directly, with local_ip and source both None.

A transport that fails before yielding anything ends serve_socket with a debug log (inbound transport failed: …); the client never reaches a core.

Deeper: TCP and TLS transports and WebSocket and gRPC transports.

serve_connection builds the connector first, since every core dials through it:

let connector = AppConnector {
router,
ctx: ctx.clone(),
};

It then matches on StreamInbound.protocol, a StreamProtocol built once per generation, and constructs a fresh core for this connection. The core borrows the inbound’s shared user table through an Arc, and the runtime’s buffer size is the core’s own BUF_SIZE constant:

StreamProtocol Driver Core constructor BUF_SIZE
Socks(SocksInbound<()>) SocksInbound::serve none: bespoke driver RELAY_BUF, 16 KiB per direction of a BidirectionalConnection
Http(..) drive HttpCore::new(config, sniff, source) MAX_HEAD, 64 KiB
Trojan(..) drive TrojanCore::new(validator, sniff, source) 16 KiB
Vless(..) drive VlessCore::new(validator, sniff, source) 16 KiB
Vmess(..) drive VMessCore::new(validator, now_unix, sniff, source) 32 KiB
Shadowsocks(..) drive ShadowsocksCore::new(resolved, sniff, source) 20 KiB
Ss2022 { config, validator } drive Ss2022Core::with_system_clock(config, validator, sniff, source) 32 KiB

The clock is a constructor argument (now_unix for VMess, the system clock for Shadowsocks 2022), never something the runtime injects, so a core stays deterministic under test.

SOCKS is the one server that does not fit ProxyCoreDecode: its UDP side lives on a second socket that the control connection only keeps alive, and that socket hears only the control connection’s client (see How UDP differs). SocksInbound::serve runs the SOCKS 4/4a/5 handshake under its own HANDSHAKE_TIMEOUT, calls connector.connect(flow) with the same AppConnector, writes the granted reply only after the connect succeeds (or a refusal code if it fails; a sniffed IP target is the exception, see step 5), and relays with a BidirectionalConnection guarded by RELAY_IDLE_TIMEOUT. Everything from step 4 on describes the core path.

Deeper: Serving and, for the SOCKS driver, SOCKS.

drive builds the runtime in stream mode:

let mut runtime =
ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress();

showing_progress turns the runtime from a Future into a Stream of Traffic deltas, one per unit of work. The app does not account traffic; it uses stream mode so it can check the core between steps.

A runtime delivers no event before the client sends its first byte, and the core’s own Timing arms its handshake deadline only on that first byte event. A client that connects and says nothing would otherwise never be timed out, so drive is the watchdog: while runtime.core().is_established() is false, every runtime.next() is wrapped in tokio::time::timeout(HANDSHAKE_TIMEOUT, …). On expiry drive returns inbound handshake timed out after 10s and the connection is dropped.

The two deadlines cover each other. drive catches a client that never speaks; once the first byte arrives, Timing::touch arms HANDSHAKE_TIMEOUT exactly once (it does not re-arm while the phase is still Handshake), so a client that trickles its request byte by byte still hits the core’s deadline 10 s after its first byte, and the core fails with client did not complete its request in time.

is_established() comes from Timing::is_established: true once the core is in Phase::Relay or Phase::Closing. A core enters Relay when it accepts the request: for a plain TCP request in the same call that pushes Open, before the dial completes; for a UDP or mux request as soon as the header is authenticated; for a sniffed request only when sniffing ends. At that point drive drops its clone of the handshake permit (handshake.take()) and polls the runtime until it ends; from here on the core’s own deadlines govern the connection.

Deeper: Server runtime and Protocol foundations.

Transport reads land in the runtime’s read buffer (up, a ReadBuffer<BUF_SIZE>). The runtime hands the core the whole unparsed region as Event::Transport(&mut [u8]); the core returns how many leading bytes it consumed. A core that needs more returns Ok(0), and the runtime reads again, keeping the unconsumed tail. If the unparsed region fills the whole buffer and the core still consumes nothing, the runtime fails with RuntimeError::FrameTooLarge.

TrojanCore is the plainest example. On its first byte event Timing::touch pushes Effect::SetDeadline(Some(HANDSHAKE_TIMEOUT)). Once parse_request_header has a whole header, the core looks the password hash up in its Validator, builds the flow and opens it:

let flow = Flow::new(header.destination, user, self.source);
// ...
fx.push(Effect::Open {
key: FlowKey::Direct,
target: flow,
});

open_tcp also enters Phase::Relay, which re-arms the deadline to RELAY_IDLE_TIMEOUT, and the call returns the header length. The runtime immediately offers the rest of the buffer again, and the core pushes a Forward { key: FlowKey::Direct, range } for the payload that arrived with the header. That forward waits: the runtime applies effects strictly in order, and a forward to a key that is still connecting stops the queue.

An unknown user makes the core return an error (trojan: invalid user, PermissionDenied), which the runtime reports as RuntimeError::Core; the connection ends with no reply.

When the inbound has sniff on and the request names an IP (worth_sniffing), the core does not open yet. It enters Phase::Sniff (deadline SNIFF_TIMEOUT, 300 ms) and copies the first payload bytes into a SniffPrefix, up to SNIFF_LIMIT (4 KiB), until a TLS SNI or an HTTP Host turns up, the budget runs out, or the deadline fires. Then it sets flow.sniffed, pushes Open, and forwards the collected bytes with Effect::ForwardHeld from its held buffer. The sniffed domain only feeds the router; the dial still goes to flow.destination.

Protocols whose client waits for a reply before it sends payload (HTTP CONNECT, SOCKS, Hysteria 2) answer “connected” first when they sniff an IP target, then collect the prefix.

Deeper: Server core, Sniffing, and each protocol’s page, for example Trojan or VLESS.

Routing runs synchronously inside the runtime’s task, at the moment the runtime applies Effect::Open. In drive_effects the runtime removes the effect, refuses a key that is already live (RuntimeError::DuplicateKey), creates the key’s KeyWaker, calls self.connector.connect(target), stores the returned future boxed in LinkState::Connecting, and schedules the key.

AppConnector::connect splits on the flow’s network:

fn connect(&mut self, flow: Flow) -> ConnectFuture {
if flow.destination.network == DialNetwork::Udp {
let link = FanOutLink::new(self.router.clone(), self.ctx.clone(), flow);
return Box::pin(async move { Ok(Outbound::Datagram(link)) });
}
let outbound = self.router.route(route_target(&flow, &self.ctx));
let fut = outbound.connect_stream(flow);
Box::pin(async move { fut.await.map(Outbound::Stream) })
}

route_target adapts the flow to the route model’s RouteTarget:

RouteTarget field Taken from Used by
remote, port flow.destination domain, cidr, geoip and port matchers
sniffed_domain flow.sniffed every domain matcher, geosite included, tries it as well as remote
inbound_tag ctx.inbound_tag inbound_tag
source flow.source.or(ctx.source) source_cidr
network flow.destination.network (Tcp or Udp) network

The flow’s own source wins over the listener’s. For a TCP inbound they are the same address; a QUIC or TUN inbound serves many peers and only the flow knows which one it came from.

Router::route calls RouteTable::pick: rules are tried in order, the first rule with any matching matcher wins, and otherwise the default outbound is returned. The result is an Arc<Outbound> shared for the whole generation.

Deeper: Route model.

Outbound::connect_stream(flow) returns a boxed StreamFuture; nothing is dialed until the runtime polls it in service_key. Outbound::Balanced calls Balancer::select() here, so a member that went down is skipped from the next flow on. Outbound::Blackhole resolves at once to OutboundStream::Blackhole, which reads EOF and swallows writes.

FreedomConnector::dial resolves flow.destination with destination_to_socketaddrs (the Resolver and AddressFamilyStrategy the outbound was built with), then TcpDialer::connect_any tries the addresses one after another, each bounded by DEFAULT_CONNECT_TIMEOUT (10 s), and returns every failure in one error (failed to connect to any address (…)). The result is OutboundStream::Tcp.

Deeper: Outbounds, Client runtime, Dialers, WireGuard, Hysteria 2 client.

8. Connected, and the reply the core stages

Section titled “8. Connected, and the reply the core stages”

service_key polls the connect future whenever the key’s waker fires. On Ready(Ok(outbound)) the slot becomes LinkState::Stream, the key is rescheduled so its first read is attempted, and the core receives Event::Connected { key }. What the core does with it is protocol-specific:

Core On Connected (TCP) On ConnectFailed
HttpCore (CONNECT) stages HTTP/1.1 200 Connection established stages HTTP/1.1 502 Bad Gateway (unless the 200 already went out for a sniffed IP target), then ShutdownTransport and Finish
VlessCore stages the 2-byte RESPONSE_HEADER closes without a reply
VMessCore stages its sealed response header closes without a reply
TrojanCore nothing: Trojan has no reply closes
ShadowsocksCore nothing closes
Ss2022Core nothing yet: the response header is sealed with the first downlink bytes, or at outbound EOF if there are none closes

Right after the core returns, handle_meta runs drive_effects, and the forwards that were waiting behind Open are written to the outbound. The staged reply goes to the client on the next poll_transport_write.

On Ready(Err(error)) the runtime calls forget_key, which removes the slot and drops every queued effect still aimed at that key, and then delivers Event::ConnectFailed { key, error }. A connect failure is an event for the core to answer, not a runtime error. Single-stream cores answer through Passthrough::on_outbound_gone, which pushes ShutdownTransport and Finish.

UDP and mux requests are answered earlier: VLESS and VMess stage their response header as soon as the request is authenticated, because those flows are opened per packet or per sub-session later.

Deeper: VLESS, VMess wire format, HTTP, Shadowsocks 2022.

From here the runtime alternates between the two directions. Each step works through, in order: the deadline, effects still waiting on an outbound, staged writes toward the client, then the client read and the ready outbounds, taking turns (prefer_transport) so neither side starves the other.

Uplink (client to target). A transport read fills up; Event::Transport hands the core the unparsed bytes; a relaying core pushes Forward { key, range } (for a plaintext core, Passthrough::on_transport forwards the whole slice). A queued forward pins the transport buffer: the runtime does not read from the client again until the forward has been written. drive_effects writes with poll_write and then poll_flush; if the outbound returns Pending, the queue waits there, and the client is not read. That is how backpressure reaches the client.

Downlink (target to client). An outbound is serviced only when its KeyWaker has fired, and read only when the staging buffer has STAGING_RESERVE + 1 bytes of room. The read goes into scratch (at most the room left above the reserve) and reaches the core as Event::Outbound { key, data }, which the core must consume whole: it stages the bytes, sealing them for an encrypted protocol. poll_transport_write writes staging to the client. A client that stops reading fills staging, and the runtime stops reading the outbound.

Idle. Timing::touch runs at the top of every byte event and re-arms RELAY_IDLE_TIMEOUT (300 s). When the deadline fires with nothing moved in either direction, Timing::expired returns Expired::Idle and pushes Finish.

When the outbound is a proxy client runtime, the core’s Forward lands in ProxyClientRuntime::poll_write, which seals the bytes with the codec into its own staging buffer and writes them to the upstream; poll_read opens frames in place and copies plaintext out. All of it runs inside the connection’s task.

Deeper: Server runtime and Links and types.

Single-stream cores keep a Passthrough for the close bookkeeping:

Event Passthrough call Effects pushed
Event::TransportEof (client finished sending) on_transport_eof Shutdown { key }, then Finish if the outbound had already ended
Event::OutboundEof { key } (target finished sending) on_outbound_eof ShutdownTransport, then Finish if the client had already ended
Event::ConnectFailed or Event::OutboundError on_outbound_gone ShutdownTransport, Finish

Effects apply in order, so Shutdown { key } reaches the outbound (poll_shutdown) only after every forward queued before it has been written. ShutdownTransport shuts the client’s write side only after staging has drained and been flushed. A slot whose read and write sides have both closed is removed from the runtime’s outbounds map, which drops (and closes) the outbound.

After Finish, the runtime stops reading from either side and only drains staging. step returns None once the core has finished, staging is empty, and any requested transport shutdown has completed. The progress stream ends, drive returns Ok(()), serve_connection returns, the runtime and everything it owns are dropped, and the connection’s clone of the session permit goes with it; the permit returns to the semaphore when the last clone for that socket is gone.

A failure ends the same way, earlier. drive turns any RuntimeError into an io::Error, and serve_connection logs it at debug level with the protocol’s StreamProtocol::name() and the Debug form of the source, for example trojan connection from Some(192.0.2.10) ended: proxy core: trojan: invalid user.

Deeper: Server core for Passthrough and the close effects, Server runtime for the drain.

A UDP flow reaches the connector the same way, as an Effect::Open whose Flow has DialNetwork::Udp, but AppConnector does not route or dial it. It returns a FanOutLink at once, so the core’s Connected arrives on the next poll and every packet is routed on its own:

flowchart TB
  core["Core: Effect::SendTo(key, to, range)"]
  send["FanOutLink::poll_send_to"]
  rt["route_target(flow.toward(to), ctx)"]
  route["Router::route: Arc of Outbound"]
  known{"Sub-link open for this outbound?"}
  write["Sub-link poll_send_to, move to most recently used"]
  busy{"Another sub-link opening?"}
  opendg["Outbound::connect_datagram(flow.toward(to))"]
  wait["Pending: the runtime keeps the SendTo queued"]
  drop["Open failed: packet dropped"]
  core --> send --> rt --> route --> known
  known -->|yes| write
  known -->|no| busy
  busy -->|yes| wait
  busy -->|no| opendg
  opendg -->|ok| write
  opendg -->|error| drop
  • Identity. Sub-links are keyed by Arc::as_ptr of the routed outbound. The router’s table holds those Arcs for the whole generation, so the pointer names “the same outbound”.
  • Routing input. Each packet routes as self.flow.toward(to): the association’s user and source, the packet’s own destination (whose network is Udp), and no sniffed domain (toward resets it).
  • One open at a time. Only one sub-link is opened at a time. A packet for a different outbound returns Pending until that open finishes, and the runtime keeps its SendTo queued. A failed open drops the packet (Ok(buf.len())) so the association does not stall on an unreachable peer.
  • Bound. At most MAX_SUBS (64) sub-links stay open; opening the 65th evicts the least recently sent-to.
  • Replies. poll_recv_from polls the sub-links round robin from where the last receive stopped and returns the first packet with its source, which reaches the core as Event::Datagram { key, from, data }. A sub-link whose receive fails is removed; the association carries on with the rest.

Outbound::connect_datagram opens the per-outbound link: a dual-stack socket for freedom, the protocol’s datagram codec over a ProxyClientRuntime for Trojan, VLESS and VMess, a SocksUdpLink for SOCKS (which takes a datagram as a reply only if endpoint of its sender equals endpoint of the relay the server named, so a dual-stack socket still hears an IPv4 relay), a tunnel association for WireGuard, a session on the shared QUIC connection for Hysteria 2. HTTP, Shadowsocks and Shadowsocks 2022 outbounds return Unsupported (… carries no datagrams), so their packets are dropped.

SocksInbound::serve runs its handshake through handshake_with_udp_source, which returns, alongside the Handshake, a UDP ASSOCIATE’s DST.ADDR and DST.PORT (None for any other request). For a UDP ASSOCIATE, associate builds an ExpectedSender from the control connection’s source, that declared source and the relay IP before it binds the relay socket (the hub):

  • One client. Over TCP, every datagram must come from the control connection’s IP, compared in canonical form by endpoint (protocols/src/socks/protocol.rs), so an IPv4-mapped IPv6 address is the IPv4 one. A datagram from any other sender is dropped before it is parsed and does not reset the association’s RELAY_IDLE_TIMEOUT.
  • Port pinning. The first datagram forwarded (one that parses and carries a payload) pins the port, and replies go to that address alone. A request that names the control connection’s own IP with a non-zero port pins that port up front. A request naming any other source is set aside, neither trusted nor refused.
  • No client to hold to. Over a Unix socket source is None, so the request must name its exact IP address and a non-zero port. A request that does not, and a hub whose family cannot receive from the client (hears: an IPv4 or IPv4-mapped hub hears only IPv4, the unspecified IPv6 hub both, any other IPv6 hub only IPv6), get the SOCKS5 reply STATUS_NOT_ALLOWED (0x02) before any socket is bound.

On the first datagram it forwards, the driver asks the same AppConnector for a UDP flow, with the control connection’s source as the flow’s source, and gets the same FanOutLink, which it drives from its own select loop.

Deeper: Outbounds, Route model, SOCKS: the association’s client, mux.cool and XUDP for UDP inside a mux carrier.

Neither inbound accepts TCP sockets, so steps 1 to 4 are replaced by the inbound’s own run loop. Both loops take a closure that, given a client address, builds the AppConnector for one runtime (the loops call it once per runtime they start), and from Effect::Open on everything above applies unchanged:

let make = move |ip| AppConnector {
router: router.clone(),
ctx: FlowContext {
inbound_tag: tag.clone(),
source: Some(ip),
},
};
Hysteria 2 (run_hysteria_inbound → Hy2Inbound::run) TUN (run_tun_inbound → TunInbound::run)
Accepts QUIC connections on one UDP socket, each under a max_connections permit (a full listener refuses the handshake) TCP connections and UDP flows terminated by a userspace IP stack on the device
Authentication once per QUIC connection, over HTTP/3 none; TunConfig.user is used for attribution
TCP one ProxyServerRuntime over Hy2StreamCore per classified proxy stream, each under a circuit_permits permit one ProxyServerRuntime over PassthroughCore per TCP connection, under a max_flows permit
UDP when the listener relays UDP, one runtime per authenticated QUIC connection over Hy2UdpCore and QuicDatagrams (over_datagrams) when UDP is enabled, one runtime per client source over TunUdpCore and TunUdpLink, under a max_flows permit, fed that source’s new flows through a FLOW_QUEUE (16) channel
On cancellation the run loop returns, then Hy2Inbound::shutdown waits for the UDP port to be released the run loop returns, then TunInbound::shutdown waits for the device to close

PassthroughCore knows its destination before the first byte, because the IP stack already knows where the client is going. It opens on the first event and relays verbatim, with optional sniffing. Since a runtime delivers no event before the client speaks, TUN’s serve_stream applies the same silent-client watchdog as drive: HANDSHAKE_TIMEOUT around each step until the core is established (tun: the client never spoke).

Deeper: Hysteria 2 server, TUN.

katana consumes the same kernel crates and keeps the same shape: accept, transport, one runtime per stream, a Connector that routes and dials. The differences are in admission, in the watchdogs around the runtime, and above all in the connector.

etemenanki-app (app/src/serve.rs) katana (src/serve.rs)
Task scope spawn_scoped(token, fut) returns a JoinHandle Scope pairs the token with a TaskTracker; Scope::shutdown waits for every task
Live connections MAX_LIVE_CONNECTIONS_PER_INBOUND (65,536) per listener, try_acquire MAX_LIVE_CONNECTIONS_PER_NODE (65,536) per listener, try_acquire; a refusal counts as a handshake failure
Protocols over streams SOCKS, HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022 VMess, VLESS, Trojan, Shadowsocks, Shadowsocks 2022, each core built from the user table current when the stream arrived
Handshake watchdog HANDSHAKE_TIMEOUT around each runtime.next() until established one HANDSHAKE_TIMEOUT around the whole loop until established; the pre-auth permit is dropped then
After the handshake poll the runtime until it ends; the core’s own deadlines apply select! over the runtime, the user’s retirement (until_retired) and PROGRESS_WATCHDOG (RELAY_IDLE_TIMEOUT + 60 s = 360 s), which is reset whenever a Traffic delta is non-zero

katana’s accept loop also bounds sockets in their transport handshake, with MAX_TRANSPORT_STAGES_PER_NODE (2,048) places that it waits for instead of refusing the socket. For every stream the transport yields, the sink calls ProxyManager::accept_stream, which takes a place from the node’s pre-auth semaphore (MAX_PREAUTH_STREAMS_PER_NODE, 512) without waiting, counts a refusal as a handshake failure, and spawns serve_stream with a clone of the socket’s live-connection permit. drive drops the pre-auth permit once the core is established. Every stream’s core is built from the user table that is current when the stream arrives (protocol.load_full()), so a user refresh never changes the table under a handshake in progress.

Every flow katana opens, whether a TCP request, a mux sub-flow, a UDP association or a Hysteria 2 stream, goes through KatanaConnector::connect:

impl Connector<Flow<UserTag>> for KatanaConnector {
type Stream = Metered<OutboundStream>;
type Datagram = FanOut;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;
}
flowchart TB
  flow["Flow of UserTag"]
  admit{"Admission::admit(tag)"}
  lease["Publish the user's lease to the connection's LeaseSlot"]
  gate["Gate::new(counter, lease)"]
  udp{"Network is UDP?"}
  fanout["FanOut: route, audit and bill each packet"]
  route["Router::route(route_target(dest, sniffed, source))"]
  check{"Outbound::Block or an audit hit?"}
  dial["connect_stream(dest)"]
  metered["Metered::new(stream, gate)"]
  refused["refused(): PermissionDenied"]
  flow --> admit
  admit -->|unknown user| refused
  admit -->|admitted| lease --> gate --> udp
  udp -->|yes| fanout
  udp -->|no| route --> check
  check -->|yes| refused
  check -->|no| dial --> metered
  1. Admit. Admission::admit looks the flow’s UserTag up in the node’s registry and returns the user’s UserCounter and lease (a CancellationToken), or None for a user the panel no longer lists or a credential now bound to another uid. The lookup happens under the same lock a user refresh commits under.
  2. Lease. The first admitted flow on a connection publishes its lease into the connection’s LeaseSlot (watch::Sender<Option<CancellationToken>>), so drive can end the whole connection when the user is retired.
  3. Route. katana’s own route_target(dest, sniffed, source) (src/router.rs) builds the RouteTarget from the destination (including its network), the sniffed domain and flow.source.or(self.source). It sets no inbound tag.
  4. Audit. A route to Outbound::Block, or a destination that RuleManager::detect reports as forbidden (and records against the uid), is refused with PermissionDenied.
  5. Dial and meter. The outbound dials flow.destination, and the stream is wrapped in Metered. Every read and write first waits on Gate::poll_open, which fails with ConnectionAborted once the user is retired and sleeps while the user’s token bucket is in debt; the transfer is then billed with Gate::sent or Gate::received.

UDP goes to FanOut, katana’s counterpart of FanOutLink, which routes, audits and bills every packet: a blocked or forbidden packet is dropped unbilled, and every packet that leaves, and every reply, passes the same Gate.

For Hysteria 2, run_hysteria builds each KatanaConnector with no lease slot; a retired user’s Hysteria flows end because their Metered outbounds refuse to move.

Deeper: katana Serving, Admission, Metering and Connector and UDP.

Invariant Mechanism Pinned by
A socket is admitted only with a free live-connection place, and the place is held for as long as any task serving the socket lives sessions.try_acquire_owned() before spawning; the Arc<OwnedSemaphorePermit> rides into every serve_connection no dedicated test
Every task of a generation dies with it spawn_scoped selects on token.cancelled(); dropping the runtime drops the transport and every outbound no dedicated test for stream inbounds; app/tests/integration/e2e_hysteria_inbound.rs → a_reload_rebinds_the_udp_port covers a listener that owns its socket
One gRPC connection yields one stream per HTTP/2 stream, and an idle one is given up on serve_h2 calls sink per accepted stream; Liveness protocols/tests/pipeline/transports.rs → one_grpc_connection_carries_many_streams; protocols/tests/unit/transports/accept.rs → a_connection_that_opens_no_stream_is_given_up_on
The handshake deadline arms once, then the idle deadline re-arms on every byte event Timing::touch, Timing::enter protocols/tests/unit/core/mod.rs → timing_arms_handshake_once_then_idle_per_byte_event, timing_reports_handshake_and_sniff_expiry_to_the_core
Bytes forwarded in the same batch as Open wait for the connect, and held bytes survive until applied drive_effects stops at a Connecting slot; held ranges resolve at apply time concepts/tests/runtime.rs → held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites
A frame larger than the read buffer fails the connection instead of growing a buffer RuntimeError::FrameTooLarge concepts/tests/runtime.rs → frame_larger_than_the_buffer_is_an_error
The flow’s own source wins over the listener’s flow.source.or(ctx.source) in route_target app/tests/unit/router.rs → the_circuits_own_source_wins_over_the_listeners, without_one_the_listeners_source_is_used, a_circuit_source_works_with_no_listener_source
Route context reaches the matchers route_target sets inbound_tag, source and network app/tests/integration/e2e_route_context.rs → inbound_tag_selects_the_route, source_cidr_matches_the_client_address, network_separates_tcp_from_udp
A sniffed domain routes an IP-addressed flow SniffPrefix sets flow.sniffed; RouteTarget::sniffed_domain app/tests/integration/e2e_sniff.rs → an_http_host_routes_an_ip_addressed_flow, a_tls_sni_routes_an_ip_addressed_flow, turning_sniffing_off_stops_the_domain_rule_matching
Connected means the upstream really accepted the flow ProxyClientConnecting waits for poll_connected (dial, handshake, flush) concepts/tests/client.rs → upstream_dial_failure_is_connect_failed_not_connected, refused_upstream_handshake_is_connect_failed_not_connected
A reply is staged only once the target is up, except when sniffing an IP target for a protocol whose client waits reply staged in the Event::Connected arm protocols/tests/unit/vless/core.rs → tcp_request_opens_and_replies_only_once_connected, a_refused_connect_closes_without_a_reply; protocols/tests/unit/http/core.rs → connect_to_a_domain_answers_200_once_connected, connect_refused_answers_502_and_closes, connect_to_an_ip_with_sniffing_replies_early_and_holds_the_prefix; protocols/tests/unit/vmess/core.rs → tcp_request_opens_and_replies_once_connected
A connect failure is an event, not a runtime error forget_key then Event::ConnectFailed concepts/tests/runtime.rs → connect_failure_reaches_the_core_as_an_event; protocols/tests/unit/trojan/core.rs → connect_failure_ends_the_connection_without_a_reply
A stalled outbound stalls the uplink instead of buffering it a queued forward pins the transport buffer concepts/tests/runtime.rs → stalled_outbound_holds_uplink_but_not_other_downlink
A proxy outbound runs inside the connection’s task ProxyClientRuntime is the outbound’s AsyncRead/AsyncWrite concepts/tests/client.rs → server_runtime_relays_through_a_client_runtime_in_one_task
Each EOF half-closes the other side, and the connection finishes after both Passthrough, Effect::Shutdown, Effect::ShutdownTransport concepts/tests/runtime.rs → half_close_propagates_both_ways_and_finishes; protocols/tests/unit/core/mod.rs → passthrough_half_closes_each_side_and_finishes_on_the_second
UDP is routed per packet and replies merge back into one association FanOutLink app/tests/integration/e2e_udp_route.rs → one_association_routes_each_peer_separately, replies_from_several_peers_merge_back_correctly
A SOCKS UDP association hears only its control connection’s client, and a request cannot widen that ExpectedSender::admits and pin; endpoint for both ends protocols/tests/unit/socks/server.rs → only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port, an_ipv4_mapped_address_is_the_ipv4_one, a_request_naming_the_peer_pins_its_port_up_front, a_request_naming_any_other_source_is_set_aside, over_a_unix_socket_the_request_must_name_the_exact_source, a_relay_that_cannot_hear_the_client_is_refused; protocols/tests/pipeline/socks.rs → udp_association_ignores_another_ip, udp_association_ignores_another_port_once_pinned, udp_association_holds_to_the_port_the_request_names, udp_association_sets_aside_a_source_it_cannot_hold_to, udp_association_is_not_widened_by_the_request, udp_association_over_a_unix_socket_needs_its_exact_source, udp_association_refuses_a_relay_that_cannot_hear_the_client, udp_link_ignores_datagrams_not_from_the_relay, udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay; protocols/tests/unit/socks/protocol.rs → endpoint_sees_through_ipv4_mapping_and_ignores_flow_info
katana admits, audits and meters every flow KatanaConnector::connect, FanOut, Metered katana tests/unit/connector.rs → a_user_the_registry_does_not_know_is_refused, a_blocked_destination_is_refused, a_forbidden_destination_is_refused_and_recorded, an_admitted_stream_is_billed_to_its_user, udp_is_billed_after_routing_and_blocked_packets_are_free
katana ends silent, stuck and retired connections drive’s handshake timeout, PROGRESS_WATCHDOG, until_retired katana tests/unit/serve.rs → a_silent_client_is_dropped_at_the_handshake_deadline, a_connection_that_stops_moving_is_dropped, a_retired_users_connection_ends
Where What happens What the client sees
A semaphore is full at accept socket dropped, debug log the connection closes before any byte
Accept error 100 ms backoff unless ConnectionAborted or Interrupted nothing; the next accept continues
Transport handshake fails or exceeds 10 s serve_socket logs inbound transport failed TLS or upgrade failure, or a close
gRPC stream on an unknown path REFUSED_STREAM that stream is reset; the connection stays up
Silent client drive times out after HANDSHAKE_TIMEOUT close
Malformed request or unknown user the core returns an error; RuntimeError::Core close, with no protocol reply
Routing cannot fail: the default outbound always answers —
SOCKS UDP ASSOCIATE with no client the relay could hold to (a Unix-socket request without an exact address and port, or a udp_bind in a family the client cannot reach) ExpectedSender::new fails before the hub is bound; PermissionDenied SOCKS5 reply 0x02, then close
SOCKS relay datagram from anyone but the association’s client dropped before parsing, without resetting the idle timer; the association carries on nothing
Dial or upstream handshake fails Event::ConnectFailed; the core answers HTTP 502, a SOCKS refusal code, a Hysteria refusal, or a close
Outbound read or write error during relay fail_outbound, then Event::OutboundError ShutdownTransport and Finish: a clean close after staged bytes
Transport read or write error RuntimeError::Transport ends the runtime the socket is already gone
Idle for RELAY_IDLE_TIMEOUT Expired::Idle pushes Finish staged bytes drain, then close
Generation cancelled (reload or shutdown) every scoped future is dropped an abrupt close of every connection on the old generation

Cancellation is by drop throughout. Nothing in the path needs an explicit shutdown call: dropping the runtime closes the transport and every outbound, and dropping the Arc permit handles returns the places to their semaphores.

Constant Value Defined in Applies to
MAX_LIVE_CONNECTIONS_PER_INBOUND 65,536 app/src/serve.rs live sockets per stream inbound
MAX_HANDSHAKES_PER_INBOUND 2,048 app/src/serve.rs connections in their handshake per stream inbound
ACCEPT_ERROR_BACKOFF 100 ms app/src/serve.rs sleep after a non-transient accept error
TRANSPORT_HANDSHAKE_TIMEOUT 10 s protocols/src/transports/accept.rs TLS, WebSocket upgrade, HTTP/2 preface
TCP_KEEPALIVE_IDLE, …_INTERVAL, …_RETRIES 120 s, 30 s, 3 protocols/src/transports/keepalive.rs every socket InboundTransport::accept receives and every socket TransportConnector::dial opens (freedom’s sockets keep the platform default)
HANDSHAKE_TIMEOUT 10 s protocols/src/core/mod.rs drive’s watchdog and the core’s handshake deadline
SNIFF_TIMEOUT, SNIFF_LIMIT 300 ms, 4 KiB protocols/src/sniff/mod.rs the sniffing window
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs idle relay
DEFAULT_CONNECT_TIMEOUT 10 s per address environment/src/dial/tcp.rs TCP connects by freedom and TransportConnector
OPEN_STREAM_TIMEOUT 5 s protocols/src/hysteria/connection.rs opening a Hysteria 2 proxy stream
WORK_BUDGET 64 units concepts/src/runtime.rs a runtime polled as a Future (the Hysteria 2 runtimes and TUN’s UDP runtimes; drive uses stream mode) yields after this many units
INLINE_EFFECTS 4 concepts/src/core.rs effects stored inline before the list spills to the heap
MAX_SUBS 64 app/src/outbound/udp_fanout.rs open sub-links per UDP association
HTTP_BUF, SOCKS_BUF, TROJAN_BUF, VLESS_BUF 16 KiB app/src/outbound/mod.rs proxy client runtime buffers
VMESS_BUF, SS2022_BUF; SS_BUF 32 KiB; 20 KiB app/src/outbound/mod.rs proxy client runtime buffers

Memory per connection follows from the buffer sizes. A server runtime allocates exactly three BUF_SIZE arrays (read, staging, scratch), so a VMess connection costs 96 KiB of buffers and an HTTP one 192 KiB. A proxy outbound adds its client runtime’s two buffers (staging and read), for example 2 × 32 KiB for VMess. Each opened outbound adds one KeyWaker Arc and one boxed connect future.

  • concepts/tests/runtime.rs drives toy cores through the runtime: open, relay, half-close, connect failure, deadlines, held bytes and backpressure. stream_mode_reports_each_unit_of_work pins the showing_progress mode drive relies on, and deadline_event_lets_the_core_time_out and deadline_is_armed_against_the_tokio_clock pin the timer behind Effect::SetDeadline.
  • concepts/tests/client.rs covers ProxyClientRuntime and the Connected semantics of ProxyClientConnector.
  • protocols/tests/unit/*/core.rs hand-drive each core, most of them through CoreHarness: when it opens, what it stages on Connected and ConnectFailed.
  • protocols/tests/pipeline/ runs each server core in a real runtime against its client codec over loopback sockets (support::pipeline::serve_runtime, client); transports.rs covers every inbound transport, including many streams on one gRPC connection.
  • app/tests/integration/ starts the real etemenanki-app binary: e2e_route_context.rs, e2e_sniff.rs and e2e_udp_route.rs exercise steps 6 and 7 and the UDP fan-out end to end.
  • app/tests/integration/e2e_unix.rs (socks_over_a_unix_socket_relays_and_cleans_up, http_connect_over_a_unix_socket_relays) covers the Unix path that skips the transport; e2e_hysteria_inbound.rs and e2e_tun.rs cover the two inbounds with their own run loops.
  • app/tests/unit/serve.rs → accept_error_backoff_classification pins the accept-error classification; it is the only unit test of app/src/serve.rs, so the app’s drive watchdog has no dedicated test. katana’s tests/unit/serve.rs → a_silent_client_is_dropped_at_the_handshake_deadline pins katana’s own watchdog, which differs in shape (one timeout around the whole handshake instead of one per step).

See Testing for how to run each layer.