Skip to content

Client codecs and the client runtime

Source files: 31 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/http/codec.rs
  • Etemenanki/protocols/src/socks/codec.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/trojan/codec.rs
  • Etemenanki/protocols/src/vless/codec.rs
  • Etemenanki/protocols/src/vmess/codec.rs
  • Etemenanki/protocols/src/ss_legacy/codec.rs
  • Etemenanki/protocols/src/ss_2022/codec.rs
  • Etemenanki/protocols/tests/support/pipeline.rs
  • Etemenanki/protocols/tests/pipeline/http.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/protocols/tests/pipeline/trojan.rs
  • Etemenanki/protocols/tests/pipeline/vless.rs
  • Etemenanki/protocols/tests/pipeline/vmess.rs
  • Etemenanki/protocols/tests/pipeline/shadowsocks.rs
  • Etemenanki/protocols/tests/unit/socks/codec.rs
  • Etemenanki/protocols/tests/unit/socks/protocol.rs
  • Etemenanki/protocols/tests/unit/vless/codec.rs
  • Etemenanki/protocols/tests/unit/vmess/codec.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • katana/src/outbound/mod.rs
  • katana/src/outbound/proxy.rs

A server core parses what a client sends and decides where each flow goes. When a flow goes to another proxy server instead of straight to its destination, something has to speak that proxy’s protocol on the way out: write the request header, wait for a reply if the protocol has one, seal every payload byte into wire frames, and open the frames that come back. In Etemenanki that is the client side of the model.

The client side has three layers, all in etemenanki-concepts:

  • a codec, a sans-I/O state machine for one flow, written once per protocol in etemenanki-protocols;
  • ProxyClientRuntime, which drives one codec over one dialed upstream and looks like a plain byte stream (or a datagram link) from outside;
  • ProxyClientConnector, a Connector that builds a codec and a runtime for every flow, so a server runtime can dial through an upstream proxy.

This page is for contributors who change any of the three or add a client protocol. It assumes you know the server side from The server core and The server runtime.

Layer Owns Does not own
Codec (ProxyCoreEncodeHandshake plus ProxyCoreEncode or ProxyCoreEncodeDatagram) The protocol: request header, handshake replies, framing, encryption, the close. Sockets, buffers, wakers, timers. It sees only byte slices and a Staging area.
ProxyClientRuntime (concepts/src/client.rs) The upstream stream, a staging buffer toward it, a read buffer from it, the handshake loop, AsyncRead/AsyncWrite or DatagramLink on the plaintext side, and checks that the codec keeps its contract. Dialing (it calls an inner Connector), choosing the codec, deciding what “connected” means to a server core.
ProxyClientConnector / ProxyClientConnecting (concepts/src/client.rs) Picking a stream or datagram codec per flow, and resolving the dial only once the upstream is dialed, the handshake is done and the handshake bytes are flushed. Anything after the dial resolves: from then on the server runtime polls the ProxyClientRuntime directly.

Nothing on the client side spawns a task or uses a channel. A ProxyClientRuntime is polled by whoever holds it, usually a ProxyServerRuntime that stores it as one of its outbounds, so the codec, both buffers and the upstream stream are all polled inside the server runtime’s task. The one exception sits below the client runtime: the gRPC client transport runs its HTTP/2 connection driver as a spawned task, which is aborted when the stream is dropped.

Every client codec implements the handshake trait. The payload shape is a second trait on top of it, so a codec implements exactly the operations its flow has.

concepts/src/core.rs
pub trait ProxyCoreEncodeHandshake {
type Target;
type Error;
const STAGING_RESERVE: usize;
fn start(&mut self, out: &mut Staging<'_>) -> Result<Handshake, Self::Error>;
fn reply(&mut self, wire: &mut [u8], out: &mut Staging<'_>) -> Result<Reply, Self::Error>;
fn finish(&mut self, out: &mut Staging<'_>) -> Result<(), Self::Error>;
}
Item Meaning
Target What the runtime’s inner Connector dials: the upstream proxy server, not the flow’s destination. The destination lives inside the codec and start encodes it. Every codec in etemenanki-protocols uses Destination.
Error The codec’s error. The runtime requires Error + Send + Sync + 'static and wraps it into an io::Error of kind InvalidData.
STAGING_RESERVE The most that one start, reply, finish or seal call stages beyond the plaintext it takes: header, length prefix, tag, padding. The runtime offers a stream seal at least STAGING_RESERVE + 1 bytes of room and a packet seal STAGING_RESERVE + plain.len(), so a codec that keeps to its bound never runs out of room.
start Called once, right after the dial succeeds. It stages the request header and says whether a reply must be parsed before plaintext may flow.
reply Called only while the handshake is AwaitReply, with the unparsed upstream bytes (mutable, so the codec can decrypt in place). It may stage the next round. Bytes left after the final reply are the first wire frames.
finish The plaintext side half-closed. The codec stages whatever tells the upstream so: a closing frame, or nothing when the wire’s own EOF carries it.

Two small enums carry the handshake state:

concepts/src/core.rs
pub enum Handshake {
Done,
AwaitReply,
}
pub enum Reply {
NeedMore,
Step { consumed: usize, next: Handshake },
}

Handshake::Done from start means the protocol sends its header and goes (VMess, VLESS, Trojan, Shadowsocks). Handshake::AwaitReply means the upstream must answer first (SOCKS method selection, an HTTP CONNECT status line). Reply::NeedMore asks for more wire bytes. Reply::Step reports that the first consumed bytes held a complete reply, and next says whether another round follows.

A Step must make progress: consume bytes, stage bytes, or end the handshake. The runtime rejects a step that does none of these (see The handshake loop).

concepts/src/core.rs
pub trait ProxyCoreEncode: ProxyCoreEncodeHandshake {
fn seal(&mut self, plain: &[u8], out: &mut Staging<'_>) -> Result<usize, Self::Error>;
fn open(&mut self, wire: &mut [u8]) -> Result<Opened, Self::Error>;
}
pub enum Opened {
Frame { consumed: usize, plain: Range<usize> },
NeedMore,
End { consumed: usize },
}
  • seal takes a prefix of plain and stages it as wire frames. It returns how many plaintext bytes it took. That may be fewer than offered, for example one frame’s worth, but must be at least one.
  • open looks at the unparsed wire bytes and opens the leading frame in place. Frame { consumed, plain } says the first consumed bytes are done with, and plain, a sub-range of them, now holds the decrypted plaintext. plain is empty for a control frame. End { consumed } means the upstream ended the flow cleanly with a closing frame. NeedMore asks for more bytes. A Frame or End must consume at least one byte and at most the slice.

A frame with empty plaintext is how a codec consumes a response header without a reply round. The VLESS response header, the VMess response header, the legacy Shadowsocks response salt, and the Shadowsocks 2022 response salt and fixed response header all open as empty frames, and the runtime skips them.

A datagram codec carries packets with a per-packet peer address over a stream to the upstream (Trojan, VLESS and VMess UDP).

concepts/src/core.rs
pub type OpenedFrom = (Opened, Option<Destination>);
pub trait ProxyCoreEncodeDatagram: ProxyCoreEncodeHandshake {
fn seal_to(
&mut self,
plain: &[u8],
to: &Destination,
out: &mut Staging<'_>,
) -> Result<Option<()>, Self::Error>;
fn open_from(&mut self, wire: &mut [u8]) -> Result<OpenedFrom, Self::Error>;
}

One seal_to is one packet, taken whole, or Ok(None) with nothing staged if it does not fit the room offered. One data frame from open_from is one packet, returned with its source. A frame returned with None as its source is not a packet (the VLESS or VMess response header, or an empty VMess chunk), and the runtime consumes and skips it. open_from follows the same progress rule as open.

How the peer address is used depends on the protocol. TrojanDatagram writes to into every packet and reads the source from every packet. VlessDatagram and VMessDatagram fix the target in the request header, ignore to, and report that target as the source of every packet they open.

concepts/src/core.rs
pub enum NoCodec<Target, Error> {
Never(Infallible, PhantomData<(Target, Error)>),
}

NoCodec is uninhabited and implements all three traits with STAGING_RESERVE = 0. A ProxyClientConnector always has both a stream codec type and a datagram codec type, so a protocol that carries only one kind names NoCodec for the other. The app and katana both alias it as NoUdp = NoCodec<Destination, io::Error> and use it for HTTP, SOCKS CONNECT and both Shadowsocks variants.

Codec File start returns finish stages STAGING_RESERVE
HttpConnect protocols/src/http/codec.rs AwaitReply (one round: status 200, anything else is an error) nothing HttpConnect::REQUEST_MAX (1024)
SocksConnect protocols/src/socks/codec.rs AwaitReply (method, optional username/password, request) nothing 528
TrojanStream, TrojanDatagram protocols/src/trojan/codec.rs Done nothing RESERVE (REQUEST_HEADER_MAX)
VlessStream, VlessDatagram protocols/src/vless/codec.rs Done nothing RESERVE (REQUEST_HEADER_MAX)
VMessStream, VMessDatagram protocols/src/vmess/codec.rs Done a terminator chunk HEADER_MAX.next_multiple_of(64)
SsStream protocols/src/ss_legacy/codec.rs Done nothing 32 + CHUNK_OVERHEAD + AddressCodec::MAX_LEN + 61
Ss2022Stream protocols/src/ss_2022/codec.rs Done nothing 2048

Codecs whose start returns Done answer reply with an error, because the runtime never calls it for them. Each protocol’s page has its wire format.

concepts/src/client.rs
pub struct ProxyClientRuntime<const BUF_SIZE: usize, Codec, Conn>
where
Codec: ProxyCoreEncodeHandshake,
Conn: Connector<Codec::Target>,
{ /* private fields */ }
impl<const BUF_SIZE: usize, Codec, Conn> ProxyClientRuntime<BUF_SIZE, Codec, Conn>
where
Codec: ProxyCoreEncodeHandshake,
Codec::Error: Error + Send + Sync + 'static,
Conn: Connector<Codec::Target>,
{
pub fn new(codec: Codec, connector: &mut Conn, target: Codec::Target) -> Self;
pub fn codec(&self) -> &Codec;
}

new asserts BUF_SIZE > Codec::STAGING_RESERVE and otherwise panics with BUF_SIZE must exceed the codec's STAGING_RESERVE or nothing can be sealed. It then calls connector.connect(target) and pins the returned future in a Box. Nothing is polled yet. The dial future first runs on the first poll of any plaintext-side method, or of the ProxyClientConnecting that wraps the runtime.

Which traits the runtime implements depends on the codec:

Codec bound Trait implemented Used as
Codec: ProxyCoreEncode AsyncRead + AsyncWrite a stream outbound
Codec: ProxyCoreEncodeDatagram DatagramLink<Addr = Destination> a datagram outbound
always Unpin required by DatagramLink: Unpin and by the Connector::Stream: AsyncRead + AsyncWrite + Unpin bound

Unpin is implemented unconditionally. That is sound because the dial future is kept as Pin<Box<_>>, the upstream stream is Unpin by the Connector::Stream bound, and the codec is never pinned.

The private state:

Field Type Meaning
codec Codec The protocol state machine.
wire Wire<Conn::Future, Conn::Stream> Connecting(Pin<Box<F>>), Up(S), or Down. Down is terminal: every later call that needs the wire fails with NotConnected (an empty poll_write still returns Ok(0)).
staging WriteBuffer<BUF_SIZE> Sealed bytes not yet written to the upstream.
down ReadBuffer<BUF_SIZE> Bytes read from the upstream. The unparsed region is what the codec has not opened yet.
opened Option<(usize, Range<usize>)> A frame opened but not fully handed to the reader: its wire length and the plaintext still to hand out, as absolute offsets into down.
wire_read_closed bool The upstream’s read side returned EOF.
ended bool The codec reported Opened::End.
finished bool finish has been staged. It is staged once only.
handshaking bool start or the last reply asked for a reply that has not been parsed yet.

ProxyClientConnector and ProxyClientConnecting

Section titled “ProxyClientConnector and ProxyClientConnecting”
concepts/src/client.rs
pub struct ProxyClientConnector<const BUF_SIZE: usize, Make, Conn, Upstream> {
make: Make,
dialer: Conn,
upstream: Upstream,
}
impl<const BUF_SIZE: usize, Make, Conn, Upstream>
ProxyClientConnector<BUF_SIZE, Make, Conn, Upstream>
{
pub fn new(make: Make, dialer: Conn, upstream: Upstream) -> Self;
}
impl<const BUF_SIZE: usize, Make, Conn, Upstream, Flow, S, D> Connector<Flow>
for ProxyClientConnector<BUF_SIZE, Make, Conn, Upstream>
where
Make: FnMut(Flow) -> Outbound<S, D>,
S: ProxyCoreEncode<Target = Upstream>,
S::Error: Error + Send + Sync + 'static,
D: ProxyCoreEncodeDatagram<Target = Upstream>,
D::Error: Error + Send + Sync + 'static,
Conn: Connector<Upstream>,
Upstream: Clone,
{
type Stream = ProxyClientRuntime<BUF_SIZE, S, Conn>;
type Datagram = ProxyClientRuntime<BUF_SIZE, D, Conn>;
type Future = ProxyClientConnecting<BUF_SIZE, S, D, Conn>;
fn connect(&mut self, flow: Flow) -> Self::Future;
}
pub struct ProxyClientConnecting<const BUF_SIZE: usize, S, D, Conn>
where
S: ProxyCoreEncodeHandshake,
D: ProxyCoreEncodeHandshake<Target = S::Target>,
Conn: Connector<S::Target>,
{
runtime: Option<
Outbound<ProxyClientRuntime<BUF_SIZE, S, Conn>, ProxyClientRuntime<BUF_SIZE, D, Conn>>,
>,
}

Flow is whatever the server side routes by: Flow in the app, Destination in katana, and String in the concept tests. Upstream is the upstream proxy’s address, cloned for every dial. ProxyClientConnecting implements Future with Output = io::Result<Outbound<ProxyClientRuntime<BUF_SIZE, S, Conn>, ProxyClientRuntime<BUF_SIZE, D, Conn>>>.

A client runtime allocates exactly two boxed [u8; BUF_SIZE] arrays (through boxed_array, so a large BUF_SIZE never touches the stack) and one boxed dial future. The buffers never grow, so backpressure comes from them being full.

flowchart LR
  W["poll_write(plain)"] --> S["codec.seal"]
  S --> ST["staging: WriteBuffer"]
  ST --> U["upstream stream"]
  U --> D["down: ReadBuffer"]
  D --> O["codec.open, in place"]
  O --> R["poll_read copies plain out"]
  • staging is FIFO. start, reply, seal, seal_to and finish all append through WriteBuffer::staging(). drain writes pending() to the upstream and advances start. When the tail is full, WriteBuffer::room() slides the pending bytes to the front, which is the only copy the buffer makes on its own. The request header staged by start therefore always leaves before the first sealed frame.
  • down has three regions: bytes already opened, [start, end) unparsed, and free space. fill reads into the free tail and, only when the tail is empty, calls compact() to move the unparsed bytes to the front. If the unparsed region fills the whole array (is_saturated()), the frame cannot fit and fill fails.
  • In place. open and open_from decrypt inside down. The runtime copies plaintext straight from down into the caller’s ReadBuf, with no intermediate buffer. A large frame read through a small ReadBuf is handed out across several calls: opened remembers the absolute range still to deliver, and down.advance_start(consumed) runs only once the last byte has gone.
stateDiagram-v2
  [*] --> Connecting: new
  Connecting --> Down: dial error or datagram link
  Connecting --> Down: start error
  Connecting --> Handshaking: start returns AwaitReply
  Connecting --> Open: start returns Done
  Handshaking --> Handshaking: Step, next AwaitReply
  Handshaking --> Open: Step, next Done
  Handshaking --> Down: EOF or wire I/O error
  Open --> Ended: codec opens End
  Open --> Down: wire I/O error
  Ended --> [*]
  Down --> [*]

Connecting, Up and Down are the Wire variants. Handshaking is Up with handshaking = true, Open is Up with the handshake done, and Ended is Open with ended = true. is_ready() is true exactly in Open and Ended. Before that, every plaintext operation returns Pending while the dial or handshake is still in progress.

ready() drives the machine from any plaintext-side poll: wire() finishes the dial and runs start, then handshake() runs reply rounds. Because both poll_write and poll_read call it, a caller that only writes still gets the handshake replies read. plaintext_waits_for_every_handshake_reply in concepts/tests/client.rs pins that a write-only caller does not deadlock.

handshake() loops while handshaking is true:

  1. drain everything staged (the greeting from start, or the round the last reply staged). It returns Pending until the upstream has accepted all of it.
  2. If down holds unparsed bytes, call codec.reply(down.unparsed(), staging):
    • Step { consumed, next } with consumed larger than the bytes given fails with codec consumed more reply than it was given.
    • Step with consumed == 0, nothing staged (measured by comparing staging.room() before and after the call) and next == AwaitReply fails with codec handshake step made no progress. Without this check the loop would spin.
    • Otherwise the runtime calls down.advance_start(consumed), sets handshaking = (next == AwaitReply) and loops.
    • NeedMore falls through to step 3.
  3. fill one chunk from the upstream. EOF here moves the wire to Down and fails with UnexpectedEof, upstream closed during the handshake.

Bytes that arrive in the same read as the final reply stay in down and are opened as the first frames. plaintext_waits_for_every_handshake_reply sends the connect reply and a data frame in one write and reads the frame back.

poll_write(plain):

  1. An empty plain returns Ok(0) at once, without driving the dial.
  2. ready(): dial, start, reply rounds.
  3. make_room(STAGING_RESERVE + 1): write staged bytes to the upstream until that much room is free. Pending here is the backpressure path: the upstream is not accepting bytes, so the caller’s write stays pending.
  4. Offer min(plain.len(), room - STAGING_RESERVE) bytes to seal. A return of 0, or more than was offered, fails with codec sealed an impossible byte count.
  5. Try to drain. A Pending drain is ignored, because the bytes are already staged and the next poll or poll_flush sends them. An error is returned.
  6. Return Ok(taken).

poll_flush runs ready(), drains everything staged and then flushes the upstream stream. poll_shutdown stages finish once (after make_room(STAGING_RESERVE)), drains, flushes and finally shuts down the upstream stream’s write side. handshake_then_seal_and_open_over_the_wire checks that a plaintext shutdown puts the close frame on the wire and then half-closes it.

The server runtime follows every poll_write that takes at least one byte of a forward with a poll_flush. If the flush is Pending, it marks the slot need_flush and flushes again the next time it services that key, so a sealed frame does not sit in staging.

poll_read(buf) first runs the same prologue as a flush (ready() then drain), with one difference: a Pending drain does not block the read. The comment in the code puts it as “a stalled wire write must not block reading”. If the handshake is not done, poll_read returns Pending. Otherwise it loops:

  1. If a frame is opened, copy as much of its plaintext as buf holds and return. Once the whole frame has been handed out, clear opened and advance down past the frame.
  2. If ended, return Ok(()) with nothing filled, which is EOF.
  3. If down has unparsed bytes, call codec.open:
    • Frame passes check_frame. An empty plain is skipped (the frame is consumed and the loop continues). Otherwise opened records the frame.
    • End passes check_frame with an empty range, is consumed, and sets ended.
    • NeedMore falls through.
  4. fill. At upstream EOF, return EOF if down is empty, or UnexpectedEof, upstream closed inside a frame, if a partial frame remains.

check_frame(len, consumed, plain) rejects consumed == 0, consumed > len, plain.start > plain.end and plain.end > consumed with codec opened an impossible frame. Each of these would otherwise make the loop spin or read outside the frame.

How the stream ends:

Upstream does poll_read returns
Sends a frame the codec opens as End EOF. Later reads return EOF too, and nothing after the closing frame is opened.
Closes the connection between frames EOF. For a codec whose finish stages nothing, this is the normal end of a flow.
Closes the connection inside a frame UnexpectedEof, upstream closed inside a frame
Closes the connection during the handshake UnexpectedEof, upstream closed during the handshake, and the wire goes Down

With a ProxyCoreEncodeDatagram codec the runtime is a DatagramLink addressed by Destination. The packets still travel over the one upstream stream.

  • poll_send_to(plain, to) runs ready(), then make_room(STAGING_RESERVE + plain.len()). If staging is already empty and there is still not enough room, the packet can never fit, and the call fails with InvalidInput, frame larger than the client runtime's buffer. If the codec answers None even though that room was free, the call fails with InvalidInput, codec refused a packet that fits its reserve. On success it tries to drain and returns Ok(plain.len()), since a packet is taken whole or not at all.
  • poll_recv_from(buf) returns one packet per call. A data frame whose source is Some is copied into buf, truncated if buf is smaller, as UDP does. The whole frame is consumed either way. A frame whose source is None is skipped. After Opened::End, or at upstream EOF with an empty down, the call fails with UnexpectedEof. A datagram link has no “end of stream”, so a server core sees the end as an outbound error for the key.

In the server runtime a failed send becomes Event::SendFailed and the key stays live, while a failed receive becomes Event::OutboundError and removes the key.

connect(flow) does three things synchronously:

  1. Clones upstream.

  2. Calls make(flow). The closure returns Outbound::Stream(codec) or Outbound::Datagram(codec), which decides the flow’s payload shape.

  3. Wraps the codec in ProxyClientRuntime::new(codec, &mut dialer, upstream), which creates the dial future, and returns it inside a ProxyClientConnecting.

ProxyClientConnecting is a future. Each poll calls the runtime’s private poll_connected: ready() (dial, start, every reply round), then drain() (every handshake byte written), then poll_flush on the upstream stream. Only when all three are done does the future resolve to Ok(Outbound::Stream(runtime)) or Ok(Outbound::Datagram(runtime)). On an error it drops the runtime, and with it the upstream connection, and resolves to that error. A flush error also moves the wire to Down before the runtime is dropped. Polling the future again after it has resolved panics with polled after completion.

This timing is what makes a server core’s events mean what they say. The server runtime turns the connector future’s result into Event::Connected or Event::ConnectFailed (see The server runtime). So:

  • Connected is delivered only after the upstream is dialed, the codec’s handshake has finished, and the handshake bytes have been written and flushed into the upstream stream. A SOCKS or HTTP server core that answers its own client on Connected therefore answers only once the upstream proxy has agreed.
  • ConnectFailed covers every failure before that point: the dial itself, a refused SOCKS method or request, an HTTP status other than 200, or an upstream that closes during the handshake. The core can send its client a real refusal instead of a premature success.

While the slot is still connecting, drive_effects stops applying the server runtime’s effect queue at the first forward, send or shutdown aimed at that key (LinkState::Connecting(_) => break); effects behind it wait too. Forwards the core pushed together with its Open therefore wait until the connecting future resolves, and the request header is always flushed on its own before the first payload frame is sealed.

Both programs wrap each proxy outbound in a ProxyClient<const BUF: usize, S, D> (app/src/outbound/proxy.rs and katana’s src/outbound/proxy.rs):

app/src/outbound/proxy.rs
pub type NoUdp = NoCodec<Destination, io::Error>;
pub type Make<S, D> = Box<dyn FnMut(Flow) -> link::Outbound<S, D> + Send>;
pub struct ProxyClient<const BUF: usize, S, D> {
inner: Mutex<ProxyClientConnector<BUF, Make<S, D>, TransportConnector, Destination>>,
}
impl<const BUF: usize, S, D> ProxyClient<BUF, S, D>
where
S: ProxyCoreEncode<Target = Destination, Error = io::Error>,
D: ProxyCoreEncodeDatagram<Target = Destination, Error = io::Error>,
{
pub fn new(make: Make<S, D>, transport: TransportConnector, server: Destination) -> Self;
pub fn connect(&self, flow: Flow) -> ProxyClientConnecting<BUF, S, D, TransportConnector>;
}
  • The inner connector is TransportConnector (protocols/src/transports/connect.rs). It dials the upstream’s Destination over TCP, TLS, WebSocket or gRPC. Its Datagram type is NoDatagram: it always yields a byte stream, because a proxy’s UDP rides inside that stream.
  • The parking_lot::Mutex is held only while connect builds the codec and the dial future. The future is awaited outside the lock.
  • The Make closure picks the codec from the flow. For Trojan, VLESS and VMess it returns a datagram codec when flow.destination.network == DialNetwork::Udp and a stream codec otherwise. HTTP, SOCKS CONNECT and Shadowsocks always return a stream codec. For a UDP flow, the HTTP and Shadowsocks outbounds return an error from connect_datagram without dialing. SOCKS UDP does not go through the client runtime: SocksOutbound builds its own UDP ASSOCIATE link (see SOCKS UDP: SocksUdpLink).
  • proxy_stream and proxy_datagram box the resolved runtime into OutboundStream::Proxy or OutboundDatagram. A resolved runtime of the wrong kind is an error (a TCP flow was dialed as datagrams, a UDP flow was dialed as a stream).

In katana the closure takes a Destination instead of a Flow (dest.network == DialNetwork::Udp picks the datagram codec), and there is no Trojan outbound. Everything else is the same.

SOCKS5 UDP does not fit a ProxyCoreEncodeDatagram codec. Its packets do not ride inside the stream to the upstream: they are separate UDP datagrams sent to a relay address that the server names in its reply. SocksOutbound::connect_datagram therefore dials the SOCKS server with TransportConnector::dial and hands that control stream to SocksUdpLink::associate (protocols/src/socks/udp_link.rs). The link implements DatagramLink<Addr = Destination> itself.

protocols/src/socks/udp_link.rs
impl<S: AsyncRead + AsyncWrite + Unpin> SocksUdpLink<S> {
pub async fn associate(
control: S,
auth: Option<(&str, &str)>,
bind: impl FnOnce(&SocketAddr) -> io::Result<UdpSocket>,
) -> io::Result<Self>;
pub fn relay(&self) -> SocketAddr;
}
  • Handshake. associate runs method selection, the optional RFC 1929 username/password round and the UDP ASSOCIATE request over control, then keeps control open for the association’s lifetime. A non-zero reply status fails with ConnectionRefused, server rejects request: <status>. The status is printed in decimal: Etemenanki’s own SOCKS inbound answers 0x02 when it cannot hold the association to its client, and the link reports that as server rejects request: 2.
  • No declared source. The request’s DST.ADDR and DST.PORT are all zeros (encode_request(CMD_UDP_ASSOCIATE, None)). RFC 1928 has a client send zeros when it does not know its source. The link binds its socket only after the reply names the relay, and behind NAT a local address would be the wrong one anyway. A server that checks sources, such as Etemenanki’s own SOCKS inbound, holds the association to the IP the control connection came from and to the port of the first datagram it relays. The server side is described in SOCKS.
  • Socket family. bind receives the relay’s address so it can pick the family. The app and katana both bind the unspecified address of the relay’s family.
  • Reply sources. poll_recv_from accepts a datagram only if it comes from the relay address the server named. Anything else is skipped before it is parsed, however well-formed. The comparison uses endpoint in protocols/src/socks/protocol.rs, the same helper the server uses: the IP in canonical form (an IPv4-mapped IPv6 address counts as its IPv4 address) and the port, with IPv6 flow info and scope left out. A dual-stack socket therefore still hears an IPv4 relay. katana 3.0.1 builds on etemenanki-protocols 2.0.1, whose link compares the socket address exactly; with the socket bound in the relay’s family, as katana binds it, the two comparisons agree. A datagram from the relay that does not parse as a relay packet is also skipped.
  • Control stream. Every poll_send_to and poll_recv_from first drains the control stream and discards what it reads. Once the control stream ends, every call fails with BrokenPipe, socks: the control connection closed. A read error on it is returned as it is, and every later call fails with BrokenPipe.
Test File Behaviour pinned
new_server_vs_new_client_udp protocols/tests/pipeline/socks.rs SocksUdpLink against the real SOCKS inbound, including a 1400-byte packet
udp_link_ignores_datagrams_not_from_the_relay protocols/tests/pipeline/socks.rs A well-formed reply from another socket is dropped, and the relay’s reply is delivered with the target address from its header
udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay protocols/tests/pipeline/socks.rs On a socket bound to [::], replies from an IPv4 relay arrive as IPv4-mapped addresses and are still accepted
endpoint_sees_through_ipv4_mapping_and_ignores_flow_info protocols/tests/unit/socks/protocol.rs endpoint equates [::ffff:1.2.3.4]:5 with 1.2.3.4:5, keeps ::1 and 127.0.0.1 apart, and ignores flow info

Data flow: a relay through a client runtime

Section titled “Data flow: a relay through a client runtime”

server_runtime_relays_through_a_client_runtime_in_one_task in concepts/tests/client.rs builds exactly this shape. A passthrough server core reads a 4-byte target and relays the rest. Its connector is a ProxyClientConnector over TinyCodec, a toy codec with a [0xC0][len:u8][target] header and XOR-masked [kind][len:u16][payload] frames, where kind 1 is data and kind 2 is the close. The server core answers + on Connected and - on ConnectFailed.

sequenceDiagram
  participant C as Client
  participant SR as ProxyServerRuntime
  participant Core as Server core
  participant CR as ProxyClientRuntime
  participant U as Upstream proxy
  Note over SR,CR: one task, no channels
  C->>SR: "dst1" + payload
  SR->>Core: Event::Transport
  Core-->>SR: Open(target), Forward
  SR->>CR: ProxyClientConnector::connect(flow) builds CR
  SR->>CR: poll ProxyClientConnecting
  CR->>U: header from codec.start, drained and flushed
  CR-->>SR: connecting future resolves Ok
  SR->>Core: Event::Connected
  Core-->>C: "+"
  SR->>CR: poll_write(payload), poll_flush
  CR->>U: sealed DATA frame
  U->>CR: DATA frame
  SR->>CR: poll_read (per-key waker fired)
  CR-->>SR: plaintext, opened in place
  SR->>Core: Event::Outbound
  Core-->>C: staged plaintext
  C->>SR: EOF
  SR->>Core: Event::TransportEof
  Core-->>SR: Effect::Shutdown
  SR->>CR: poll_shutdown
  CR->>U: CLOSE frame from codec.finish, then write-side shutdown
  U->>CR: CLOSE frame
  CR-->>SR: poll_read EOF (Opened::End)
  SR->>Core: Event::OutboundEof

Every poll the server runtime makes on the outbound uses the per-key waker from its wake-up registry. The client runtime passes that Context down to the upstream stream, so upstream readiness wakes exactly that key.

The test pins three things beyond the round trip:

  • The upstream pipe holds only 4 bytes and the header is 6 (\xC0\x04dst1), so the header cannot leave until the upstream reads. The client sees no + during that time.
  • The client codec encoded the server core’s target (dst1), not the upstream’s address (upstream).
  • The server runtime’s Traffic counts what crosses the outbound’s plaintext side, not the upstream wire: outbound_tx == 23 (the payload, with no header or framing) and outbound_rx == 8.
Invariant Enforced by Pinned by
No plaintext is sealed or delivered before the handshake is Done ready() at the top of every plaintext-side poll; the is_ready() guard in poll_read and poll_recv_from plaintext_waits_for_every_handshake_reply
The request header leaves before the first frame staging is FIFO; start stages the header before any seal handshake_then_seal_and_open_over_the_wire
A resolved dial means dialed, handshaken and flushed ProxyClientConnecting::poll calls poll_connected (ready, drain, poll_flush) server_runtime_relays_through_a_client_runtime_in_one_task
Upstream refusals become ConnectFailed, not Connected Dial and handshake errors resolve the connecting future with Err upstream_dial_failure_is_connect_failed_not_connected, refused_upstream_handshake_is_connect_failed_not_connected
A codec contract violation fails and never spins The no-progress and consumed > len checks in handshake; check_frame; the seal count check handshake_step_without_progress_is_rejected, zero_length_frame_is_rejected, end_marker_beyond_the_slice_is_rejected, datagram_zero_length_frame_is_rejected
Memory per flow is bounded Two fixed BUF_SIZE arrays; make_room waits for the upstream instead of growing backpressure_from_the_wire_reaches_the_writer
Upstream frames larger than the buffer are rejected fill fails when down.is_saturated() not tested directly
down is never compacted under an opened frame poll_read returns before fill while opened is set; debug_assert! in fill small_reads_take_one_frame_across_several_calls exercises the path
A failed dial stays failed Wire::Down is terminal and every later call returns NotConnected dial_failure_surfaces_on_first_use
A truncated upstream frame is an error, not a clean EOF eof() checks for unparsed bytes at wire EOF truncated_upstream_frame_is_unexpected_eof
One receive is one packet poll_recv_from returns after one data frame datagram_codec_sends_and_receives_packets_over_the_stream
BUF_SIZE leaves room to seal assert! in ProxyClientRuntime::new enforced at construction

Every failure reaches the caller as an io::Error:

Condition Kind Message Wire afterwards
Inner dial fails the dial error’s kind the dial error Down
Inner connector yields a datagram link Unsupported proxy client runtime needs a stream to the upstream, got a datagram socket Down
Any later call after Down NotConnected none Down
start fails InvalidData the codec error Down
reply, seal, open, finish, seal_to or open_from fails InvalidData the codec error unchanged
Upstream write or read fails the I/O error’s kind the I/O error Down
Upstream write returns 0 WriteZero none Down
Upstream flush fails while connecting the I/O error’s kind the I/O error Down
Upstream flush or shutdown fails in poll_flush or poll_shutdown the I/O error’s kind the I/O error unchanged
Upstream closes during the handshake UnexpectedEof upstream closed during the handshake Down
Upstream closes inside a frame UnexpectedEof upstream closed inside a frame unchanged
Datagram receive after the flow ended UnexpectedEof none unchanged
Unparsed upstream bytes fill down InvalidData upstream frame larger than the client runtime's buffer unchanged
A packet cannot fit even into empty staging InvalidInput frame larger than the client runtime's buffer unchanged
Codec refuses a packet that fits its reserve InvalidInput codec refused a packet that fits its reserve unchanged
Step consumed more than it was given Other codec consumed more reply than it was given unchanged
Step made no progress Other codec handshake step made no progress unchanged
open or open_from returned an impossible frame Other codec opened an impossible frame unchanged
seal took 0 or more than it was offered Other codec sealed an impossible byte count unchanged

A codec error always has kind InvalidData, whatever kind the codec used inside, because codec_err wraps it with io::Error::new(io::ErrorKind::InvalidData, e). The message is the codec’s own. HTTP’s proxy responded with status 407 keeps its text (new_server_refuses_bad_credentials_with_407 in protocols/tests/pipeline/http.rs), and refused_handshake_fails_the_flow checks both the InvalidData kind and the method refused text.

The server runtime treats any error from poll_write, poll_flush, poll_shutdown, poll_read or poll_recv_from as fatal for that outbound: it drops the outbound and tells the core with OutboundError. A poll_send_to error is reported as SendFailed and the key stays live. In practice the rows marked “unchanged” therefore do not leave a half-working stream outbound behind.

Cancellation. The client runtime owns no task and no channel, and holds nothing outside itself. Dropping it (because the server runtime closed the key, the connection finished, or the ProxyClientConnecting future was dropped before it resolved) drops the dial future or the upstream stream, and both buffers, at once. There is nothing to join.

Constant Where Value Meaning
BUF_SIZE ProxyClientRuntime const generic per protocol, below Size of staging and of down. The largest upstream frame the codec opens must fit down.
HTTP_BUF app/src/outbound/mod.rs, katana src/outbound/mod.rs 16 * 1024 HTTP CONNECT client
SOCKS_BUF both 16 * 1024 SOCKS CONNECT client
TROJAN_BUF app/src/outbound/mod.rs 16 * 1024 Trojan client (the app only)
VLESS_BUF both 16 * 1024 VLESS client
VMESS_BUF both 32 * 1024 VMess client
SS_BUF both 20 * 1024 Shadowsocks (AEAD) client
SS2022_BUF both 32 * 1024 Shadowsocks 2022 client
HttpConnect::REQUEST_MAX protocols/src/http/codec.rs 1024 Largest CONNECT request start stages; a longer one fails with http: CONNECT request exceeds the codec's reserve

Per flow, a client runtime allocates 2 * BUF_SIZE bytes plus the boxed dial future, on top of the server runtime’s own buffers. A VMess flow therefore holds 64 KiB of client-side buffers. The largest packet poll_send_to can accept is BUF_SIZE - STAGING_RESERVE, because an empty staging offers exactly BUF_SIZE bytes of room; a larger packet fails with frame larger than the client runtime's buffer.

The client runtime carries one flow per upstream connection. A client that shares one upstream carrier among many flows (a mux.cool client) would need an admission path for flows that the client runtime does not provide.

  1. Implement ProxyCoreEncodeHandshake with Target = Destination and Error = io::Error, the bounds ProxyClient in the app and katana requires. Encode the flow’s destination in start. Return AwaitReply only if the upstream must answer before payload, and make every Step consume, stage or finish.

  2. Declare an honest STAGING_RESERVE: the largest header, handshake round, closing frame or per-seal overhead you ever stage. Then Staging::reserve and Staging::put cannot fail for a codec that keeps to it.

  3. Implement ProxyCoreEncode (take at most one frame per seal, open one frame per open) or ProxyCoreEncodeDatagram (one packet per seal_to and per data frame). Consume a response header as an empty frame instead of adding a reply round.

  4. Pick a BUF_SIZE larger than STAGING_RESERVE and larger than the biggest frame the upstream may send, add a *_BUF constant and an outbound variant built with a Make closure in app/src/outbound/mod.rs, and do the same in katana’s src/outbound/mod.rs if katana needs it.

  5. Add a unit test that drives the codec with no I/O in protocols/tests/unit/<protocol>/codec.rs, included from the codec file with #[cfg(test)] #[path = …] mod tests; as the existing codecs do. Add a new_server_vs_new_client_* pipeline test against the real server core.

The concept tests in concepts/tests/client.rs drive the runtime with toy codecs (TinyCodec, TwoRoundCodec, Scripted, TinyUdpCodec) over tokio::io::duplex pipes, so each behaviour can be isolated. Run them with cargo test -p etemenanki-concepts --test client.

Test Behaviour pinned
codec_is_driven_without_any_io A codec is a plain state machine: start stages the header, seal takes at most one frame, open decrypts in place and returns NeedMore on a partial frame.
handshake_then_seal_and_open_over_the_wire The header precedes the first frame, long writes split into several frames, empty frames are skipped, End is EOF, and shutdown sends the close frame and then half-closes the wire.
small_reads_take_one_frame_across_several_calls One opened frame is handed out through 3-byte reads without being reopened.
truncated_upstream_frame_is_unexpected_eof Wire EOF with a partial frame in down is UnexpectedEof.
dial_failure_surfaces_on_first_use The dial error (ConnectionRefused) reaches the first write, and later calls get NotConnected.
backpressure_from_the_wire_reaches_the_writer With a 16-byte upstream pipe, a 200-byte write_all stays pending until the upstream reads.
plaintext_waits_for_every_handshake_reply A two-round handshake: nothing but the greeting before the first reply, no payload before the last one, a write-only caller does not deadlock, and data sent with the final reply is opened.
refused_handshake_fails_the_flow A codec error in reply surfaces as InvalidData carrying the codec’s message.
handshake_step_without_progress_is_rejected The no-progress check fires within a 1 s timeout instead of spinning.
zero_length_frame_is_rejected check_frame rejects consumed == 0.
end_marker_beyond_the_slice_is_rejected check_frame rejects an End that consumes more than the slice.
datagram_zero_length_frame_is_rejected The same guard on poll_recv_from.
datagram_codec_sends_and_receives_packets_over_the_stream seal_to frames a packet with its port, two packets in one wire write come back as two receives, and a close frame turns into UnexpectedEof.
server_side::server_runtime_relays_through_a_client_runtime_in_one_task The full relay in one task, Connected only after the header has left through a 4-byte pipe, Shutdown becoming finish, and plaintext Traffic counts.
server_side::upstream_dial_failure_is_connect_failed_not_connected A refused dial reaches the server core as ConnectFailed (the client sees -, never +).
server_side::refused_upstream_handshake_is_connect_failed_not_connected Nothing is answered while the upstream handshake is pending, and a refused method becomes ConnectFailed.

The real codecs run against the real server cores in protocols/tests/pipeline/*.rs, through the helper client in protocols/tests/support/pipeline.rs, which builds a ProxyClientRuntime over plain TCP. Run them with cargo test -p etemenanki-protocols --test pipeline.

Test File Behaviour pinned
new_server_vs_new_client_tcp http.rs, socks.rs, trojan.rs, vless.rs, vmess.rs A stream codec’s echo round trip of 70 000 to 100 000 bytes against the real server core, then a shutdown and a clean EOF
new_server_vs_new_client_udp trojan.rs, vless.rs, vmess.rs Datagram codecs through DatagramLink, including a 1500-byte packet (4000 bytes for VMess)
ss_new_server_vs_new_client_tcp, ss2022_new_server_vs_new_client_tcp, ss2022_multi_user_new_server_vs_new_client shadowsocks.rs Shadowsocks and Shadowsocks 2022 stream codecs
new_server_refuses_bad_credentials_with_407 http.rs An HTTP 407 fails the first write, and the message carries the status
new_server_refuses_an_unreachable_target_after_trying socks.rs A SOCKS refusal (rejects request) fails the first write, not a later read

Each codec also has unit tests under protocols/tests/unit/<protocol>/codec.rs that call it without I/O, the way codec_is_driven_without_any_io does. The codec file includes them as a #[cfg(test)] module, so they run with the library tests (cargo test -p etemenanki-protocols --lib). For example, anonymous_connect_takes_two_rounds and credentials_add_a_round_and_a_refusal_is_an_error cover SOCKS, stream_codec_consumes_the_response_header_as_an_empty_frame covers VLESS, and datagram_codec_refuses_a_packet_that_does_not_fit covers VMess.