Skip to content

Transports: WebSocket and gRPC

Source files: 22 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/protocols/src/transports/ws/endpoint.rs
  • Etemenanki/protocols/src/transports/ws/stream.rs
  • Etemenanki/protocols/src/transports/grpc/framing.rs
  • Etemenanki/protocols/src/transports/grpc/liveness.rs
  • Etemenanki/protocols/src/transports/grpc/settings.rs
  • Etemenanki/protocols/src/transports/grpc/stream.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/transports/stream.rs
  • Etemenanki/app/src/transport.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/serve.rs
  • Etemenanki/protocols/tests/unit/transports/ws_endpoint.rs
  • Etemenanki/protocols/tests/unit/transports/grpc_framing.rs
  • Etemenanki/protocols/tests/unit/transports/grpc_liveness.rs
  • Etemenanki/protocols/tests/unit/transports/accept.rs
  • Etemenanki/protocols/tests/pipeline/transports.rs
  • Etemenanki/app/tests/integration/e2e_xray.rs
  • Etemenanki/app/tests/integration/e2e_xray_vmess.rs
  • Etemenanki/app/tests/integration/e2e_xray_mux.rs
  • katana/src/inbound.rs

WebSocket and gRPC are the two carriers that make a proxy connection look like web traffic. Both live in etemenanki-protocols under protocols/src/transports/, and both end in the same place: a TransportStream variant that implements AsyncRead + AsyncWrite, so a protocol core never learns which carrier it runs over.

This page is for contributors who change the carriers themselves: the upgrade and early-data handling of WebSocket, the Hunk/MultiHunk framing of gRPC, the HTTP/2 connection that one served socket turns into, and the timers that reclaim dead peers. TCP, TLS, MaybeTlsStream and TCP keepalive are covered in Transports: TCP and TLS.

Piece File → symbol Does Leaves to others
WebSocket endpoint ws/endpoint.rs → WsRoute, WsTarget Path normalisation, ?ed= parsing, Host check, early-data encoding, tungstenite limits TLS (below it), the payload (above it)
WebSocket stream ws/stream.rs → WsStream Binary messages as bytes, Ping/Pong, Close as EOF, keepalive Ping, idle teardown, deferred client upgrade Protocol framing inside the payload
gRPC framing grpc/framing.rs → encode_hunk, encode_multi_hunk, HunkDecoder gRPC length prefix and the protobuf data field, message size cap HTTP/2 framing (the h2 crate)
gRPC settings grpc/settings.rs → GrpcPaths, GrpcMode h2 builder settings, the /<service>/Tun and /<service>/TunMulti paths, request and response headers, stream close
gRPC stream grpc/stream.rs → GrpcStream One HTTP/2 stream as bytes, writes under flow control, window credit, the client connection driver Connection-level supervision on the server
Serving accept.rs → InboundTransport::accept, serve_h2 The 10 s transport handshake bound, fanning an HTTP/2 connection out into streams Per-stream serving, which the caller’s sink does
Liveness grpc/liveness.rs → Liveness Idle deadline and PING/PONG for a served HTTP/2 connection
Dialing connect.rs → TransportKind, TransportConnector Resolve, connect, keepalive, then wrap in TLS, WebSocket or gRPC

Every accepted or dialed socket gets TCP keepalive first, then optional TLS, then the carrier:

flowchart LR
  tcp["TcpStream"] --> ka["set_keepalive"]
  ka --> tls["MaybeTlsStream"]
  tls --> ws["WsStream"]
  tls --> h2["h2 Connection"]
  h2 --> g1["GrpcStream"]
  h2 --> g2["GrpcStream"]
  ws --> ts["TransportStream"]
  g1 --> ts
  g2 --> ts
  ts --> core["protocol core"]

One accepted WebSocket socket yields exactly one stream. One accepted gRPC socket yields one stream per HTTP/2 stream the peer opens, for as long as the connection lives.

The inbound and outbound enums are how callers reach both carriers; WsStream::accept, WsStream::connect and GrpcStream::connect are public too, but the app and katana only use the enums and their helper constructors.

protocols/src/transports/accept.rs
pub const TRANSPORT_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
pub type Accepted = TransportStream;
pub enum InboundTransport {
Tcp,
Tls(ServerConfig),
Ws {
route: WsRoute,
tls: Option<ServerConfig>,
},
Grpc {
paths: Arc<GrpcPaths>,
tls: Option<ServerConfig>,
},
}
impl InboundTransport {
pub fn ws(path: impl AsRef<str>, host: Option<&str>, tls: Option<ServerConfig>) -> Self;
pub fn grpc(service: impl AsRef<str>, tls: Option<ServerConfig>) -> Self;
pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);
}
protocols/src/transports/connect.rs
pub enum TransportKind {
Tcp,
Tls(ClientConfig),
Ws {
target: WsTarget,
tls: Option<ClientConfig>,
},
Grpc {
authority: Arc<str>,
service: Arc<str>,
mode: GrpcMode,
user_agent: Option<Arc<str>>,
tls: Option<ClientConfig>,
},
}
impl TransportKind {
pub fn ws(host: impl AsRef<str>, path: impl AsRef<str>, tls: Option<ClientConfig>) -> Self;
pub fn grpc(
authority: impl AsRef<str>,
service: impl AsRef<str>,
tls: Option<ClientConfig>,
) -> Self;
pub fn multi(mut self) -> Self;
pub fn user_agent(mut self, agent: Option<&str>) -> Self;
}

InboundTransport::accept returns as soon as the one stream is handed to sink for WebSocket. For gRPC it returns only when the HTTP/2 connection ends, because the future itself is the connection driver (see Serving an HTTP/2 connection). The caller therefore runs accept in its own task and spawns one task per yielded stream; the app does this in serve_socket, described in Serving inbounds.

Setting etemenanki-app inbound etemenanki-app outbound katana inbound
WebSocket path ws.path, default "/" ws.path, default "/" node path, "/" if empty
WebSocket Host ws.host, or no check ws.host, then tls.server_name, then server; error ws stream needs ws.host or server node host, or no check if empty
WebSocket TLS ALPN Alpn::Http1 Alpn::Http1 Alpn::None
gRPC service grpc.service_name, required grpc.service_name, required node service name
gRPC :authority not checked grpc.authority, then tls.server_name, then server; error grpc stream needs grpc.authority or server not checked
gRPC TLS ALPN Alpn::Http2 Alpn::Http2 Alpn::Http2
gRPC mode and user agent serves both paths always GrpcMode::Gun, DEFAULT_USER_AGENT serves both paths

A missing grpc.service_name fails validation in app/src/transport.rs → resolve_stream with grpc stream needs grpc.service_name. TransportKind::multi and TransportKind::user_agent exist in the library, but the app does not call them. The user-facing keys are documented in Transports.

protocols/src/transports/ws/endpoint.rs
pub(crate) const MAX_EARLY_DATA: usize = 16 * 1024;
pub(crate) const MAX_WS_MESSAGE_LEN: usize = 1024 * 1024;
pub(crate) fn ws_config() -> tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
pub(crate) type EarlyDataSlot = Arc<Mutex<Option<Bytes>>>;
pub(crate) fn normalize_path(path: &str) -> String;
pub(crate) fn parse_early_data_path(path: &str) -> (String, Option<usize>);
pub(crate) fn host_matches(request_host: &str, config: &str) -> bool;
pub(crate) fn decode_early_data_header(value: &str) -> Result<Option<Bytes>, ()>;
pub(crate) fn encode_early_data_header(bytes: &[u8]) -> String;
pub struct WsRoute {
path: Arc<str>,
host: Option<Arc<str>>,
}
impl WsRoute {
pub fn new(path: impl AsRef<str>) -> Self;
pub fn host(mut self, host: impl Into<Arc<str>>) -> Self;
pub(crate) fn callback(&self, early_data: EarlyDataSlot) -> WsCallback;
fn accepts(&self, path: &str, host: Option<&str>) -> bool;
}
pub struct WsTarget {
host: Arc<str>,
path: Arc<str>,
early_data_limit: Option<usize>,
}
impl WsTarget {
pub fn new(host: impl AsRef<str>, path: impl AsRef<str>) -> Self;
pub fn early_data_limit(&self) -> Option<usize>;
pub(crate) fn client_request(
&self,
tls: bool,
early_data: Option<&Bytes>,
) -> std::io::Result<http::Request<()>>;
}

WsRoute is what an inbound accepts; WsTarget is what an outbound requests. WsCallback implements tungstenite’s Callback and runs inside the server handshake with a clone of the route and the shared EarlyDataSlot. The parking-lot Mutex in the slot is locked at most twice per connection: by the callback when it stores valid early data, and by WsStream::accept when it takes it.

protocols/src/transports/ws/stream.rs
pub const WS_IDLE_TIMEOUT: Duration = Duration::from_secs(300);
pub const WS_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(60);
pub struct WsStream<S> {
state: State<S>,
pending_read: Bytes,
read_closed: bool,
close_sent: bool,
control: Option<Message>,
keepalive: Pin<Box<Sleep>>,
idle: Pin<Box<Sleep>>,
read_waker: Option<Waker>,
}
impl<S> WsStream<S>
where
S: AsyncRead + AsyncWrite + Unpin,
{
pub fn from_upgraded(ws: WebSocketStream<S>, initial: Option<Bytes>) -> Self;
pub async fn accept(stream: S, route: &WsRoute) -> io::Result<Self>;
pub async fn connect(stream: S, target: &WsTarget, secure: bool) -> io::Result<Self>
where
S: Send + 'static;
pub fn is_open(&self) -> bool;
}
impl<S> AsyncRead for WsStream<S> where S: AsyncRead + AsyncWrite + Unpin { /* … */ }
impl<S> AsyncWrite for WsStream<S> where S: AsyncRead + AsyncWrite + Unpin + Send + 'static { /* … */ }

In TransportStream the stream is always WsStream<MaybeTlsStream>, boxed.

WsRoute::new runs the configured path through normalize_path, which strips ?ed= with parse_early_data_path before it mirrors Xray’s GetNormalizedPath (empty becomes /, a missing leading / is added). WsTarget::new calls parse_early_data_path itself to keep the limit, then normalize_path. The same path string can therefore be written on both sides: the server strips ?ed=N just as the client does, and never sees it on the wire.

Configured path Upgrade path Early-data limit (WsTarget)
"" / none
echo /echo none
/x?ed=2048 /x Some(2048)
x?ed=2048 /x Some(2048)
/x?foo=1&ed=2048 /x?foo=1 Some(2048)
/x?foo=1 /x?foo=1 none

The rules in parse_early_data_path:

  • Only a query part named ed with a non-empty value is treated as early-data configuration. It is removed from the path whether or not its value parses.
  • The value is parsed as usize. A value that does not parse removes the parameter but leaves the limit at None.
  • When ed is present, empty query parts are dropped and the remaining ones are re-joined with &. Without ed, the path is returned untouched.
  • WsStream::connect only defers the upgrade when the limit is greater than zero, so ?ed=0 means no early data.

host_matches mirrors Xray’s internet.IsValidHTTPHost: both sides are lower-cased, the request value is cut at its last : to drop a port, and the host part must equal the configured value exactly.

Route Request Result
no host configured any Host, or none accepted
example.com Example.COM:443 accepted
example.com other.example:443 404
example.com:443 example.com:443 404: the port is stripped from the request value only, never from the configured one
example.com no Host header, or a value that is not visible ASCII 404
any a path other than the route’s, compared case-sensitively 404

The rejection is built by reject_404: an empty-body 404 Not Found returned from the callback, which makes tungstenite answer and fail the handshake.

Xray-compatible early data carries the first client bytes inside the upgrade request, in the Sec-WebSocket-Protocol header, so a client-first protocol saves one round trip.

Step Side Mechanism
Encode client encode_early_data_header: base64 URL-safe alphabet, no padding (b"hello" → aGVsbG8, [0xfb, 0xff] → -_8)
Size client The first min(len, ed, MAX_EARLY_DATA) bytes of the first non-empty write, so never more than 16 KiB
Normalise server Trim whitespace, map + to - and / to _, drop =: standard and URL-safe, padded or not, all decode
Decode server URL_SAFE_NO_PAD. An empty value, a value that does not decode, or one that decodes to nothing is not early data
Cap server More than MAX_EARLY_DATA decoded bytes returns Err(()), and the callback answers 413 Payload Too Large
Echo server Valid early data is stored in the slot and the request’s Sec-WebSocket-Protocol value is echoed verbatim in the 101 response
Deliver server WsStream::accept takes the slot and seeds pending_read, so the early data is the first thing the core reads

A header that is not early data leaves the slot empty and the response without Sec-WebSocket-Protocol; the upgrade still succeeds.

The echo is load-bearing. tungstenite’s client handshake treats a request that carried Sec-WebSocket-Protocol and a 101 without one as a subprotocol error and fails the upgrade, so a server that consumed the early data but did not echo it would break every Etemenanki client that uses ?ed=.

On the client, WsStream::connect with a positive limit does not upgrade at all. It parks the socket in State::Deferred and returns. The first non-empty poll_write calls start_upgrade, which takes up to limit bytes, builds the request with them, boxes the tungstenite handshake into State::Upgrading, wakes any reader parked in read_waker, and returns Ok(take). The boxed handshake does nothing until the stream is polled again: the rest of a write_all, a woken read, a flush or a shutdown drives it through poll_open. The caller’s write_all then writes the rest as ordinary messages once the upgrade completes.

sequenceDiagram
  participant CC as client core
  participant CW as WsStream client
  participant SW as WsStream::accept
  participant SC as server core
  CC->>CW: connect with path /ws?ed=2048
  Note over CW: State Deferred, no bytes sent
  CC->>CW: write_all 2100 bytes, first poll_write
  CW->>CW: start_upgrade takes 2048, state Upgrading
  CW-->>CC: Ok(2048)
  CC->>CW: second poll_write, 52 bytes
  Note over CW: poll_open drives the boxed handshake
  CW->>SW: GET /ws with Sec-WebSocket-Protocol base64url(2048 bytes)
  SW->>SW: WsRoute::accepts path and Host
  SW->>SW: decode_early_data_header into EarlyDataSlot
  SW-->>CW: 101 with Sec-WebSocket-Protocol echoed
  Note over CW: Upgrading becomes Open
  SW->>SC: WsStream with pending_read = 2048 bytes
  CW->>SW: Binary message, 52 bytes
  SC->>SC: reads 2048 bytes, then 52

A read issued before the first write returns Pending from poll_open and stores its waker; start_upgrade wakes it. A poll_shutdown in Deferred starts the upgrade with no early data, so the close still reaches the peer as a WebSocket Close frame. poll_flush in Deferred returns Ok(()) without doing anything.

stateDiagram-v2
  [*] --> Open: accept, or connect without ed
  [*] --> Deferred: connect with ed greater than 0
  Deferred --> Upgrading: first non-empty poll_write, or poll_shutdown
  Upgrading --> Open: handshake done, timers reset
  Upgrading --> Failed: handshake error
  Failed --> Failed: every call returns the stored error
  Open --> [*]

State::Failed(io::ErrorKind, String) keeps the kind and text of the upgrade error, and poll_open rebuilds an equal io::Error on every later call. The write that started the upgrade has already returned Ok, so the failure surfaces on the next read, write, flush or shutdown. Dropping the stream in Upgrading drops the boxed handshake future and the socket with it.

start_upgrade takes the socket out of Deferred before it builds the request. If WsTarget::client_request fails there (a host that does not form a valid URI, for example), that first write returns the error, the state stays Deferred without a socket, and any later write or shutdown fails with websocket upgrade already started. A flush still returns Ok(()), and a read stays Pending, because nothing is left to wake it.

The mapping between WebSocket messages and bytes is fixed:

Direction Event WsStream action
read Binary(bytes) Stored in pending_read; handed out across as many poll_read calls as the buffer sizes need
read Ping(payload) control = Some(Pong(payload)), sent by the next poll_control
read Close(_) read_closed = true: EOF
read Text, Pong, raw Frame Ignored, but they still reset both timers
read stream end, ConnectionClosed, AlreadyClosed, ResetWithoutClosingHandshake EOF
read any other tungstenite error, including a size limit Err, io::Error::other unless it is an I/O error
write empty buffer Ok(0), no message
write non-empty buffer Exactly one Binary message with a copy of the buffer, then Ok(buf.len())
flush Sends a waiting control frame, then flushes the sink
shutdown One Close(None) (guarded by close_sent), then flush; ConnectionClosed or AlreadyClosed during that flush counts as success

A write whose sink reports ConnectionClosed or AlreadyClosed fails with BrokenPipe and the text websocket closed (ws_err).

control holds at most one control frame. poll_control is called from poll_read as well as the write path, and on the read path a sink that is not ready does not block the read: the frame waits and the read proceeds.

tungstenite 0.30 also queues a Pong of its own when it reads a Ping, and a Pong handed to its sink replaces a queued Pong that has not been flushed yet. The peer therefore sees one Pong per Ping, or two when tungstenite flushed its own before WsStream sent the second; both carry the Ping’s payload, and RFC 6455 allows unsolicited Pongs.

tungstenite sends no keepalive of its own, so WsStream owns two tokio::time::Sleep timers:

Timer Constant Fires Effect
keepalive WS_KEEPALIVE_INTERVAL = 60 s 60 s after the last activity Re-armed for another 60 s; queues Ping with an empty payload if no control frame is waiting
idle WS_IDLE_TIMEOUT = 300 s 300 s after the last activity read_closed = true: the read side reports EOF, as a peer that vanished without a FIN would

touch resets both. It runs when the upgrade completes, on every incoming frame of any kind (a Pong included, which is what turns the Ping into a liveness probe), and after every successful write. Both timers are polled from poll_read, so they fire while a reader is waiting on the stream. The values match the HTTP/2 idle and keepalive constants so the two carriers behave alike.

ws_config replaces tungstenite’s defaults (64 MiB per message, 16 MiB per frame) with MAX_WS_MESSAGE_LEN = 1 MiB for both max_message_size and max_frame_size. tungstenite reassembles a fragmented message into one buffer before yielding it, so this is the most memory a single peer can make one session hold for a message. The same config is passed to accept_hdr_async_with_config and client_async_with_config, so it applies to inbound and dialed sessions alike, and it matches the gRPC carrier’s MAX_GRPC_MESSAGE_LEN.

Both tungstenite limits apply to received messages only. WsStream::poll_write never splits a buffer, so a single write larger than 1 MiB goes out as one message that an Etemenanki peer rejects with a size error. The cores write far smaller chunks, and the pipeline tests send 200 000 bytes in one write.

The gRPC carrier is Xray’s “gun” transport: a bidirectional-streaming gRPC call whose messages are protobuf Hunk { bytes data = 1; } or, in multi mode, MultiHunk { repeated bytes data = 1; }. There is no generated protobuf code; the one field is written and parsed by hand.

protocols/src/transports/grpc/settings.rs
pub struct GrpcPaths {
tun: Arc<str>,
multi: Arc<str>,
}
impl GrpcPaths {
pub fn new(service: impl AsRef<str>) -> Self;
pub fn classify(&self, path: &str) -> Option<GrpcMode>;
}
pub enum GrpcMode {
Gun,
Multi,
}
pub(crate) fn configured_client_builder() -> h2::client::Builder;
pub(crate) fn configured_server_builder() -> h2::server::Builder;
pub(crate) fn grpc_response() -> http::Response<()>;
pub(crate) fn grpc_request(
secure: bool,
authority: &str,
service: &str,
mode: GrpcMode,
user_agent: Option<&str>,
) -> std::io::Result<http::Request<()>>;
pub(crate) fn finish_stream(send_stream: &mut SendStream<Bytes>, is_server: bool);
protocols/src/transports/grpc/framing.rs
pub(crate) fn encode_hunk(payload: &[u8]) -> io::Result<Bytes>;
pub(crate) fn encode_multi_hunk(payloads: &[Bytes]) -> io::Result<Bytes>;
pub(crate) struct HunkDecoder {
chunks: VecDeque<Bytes>,
len: usize,
}
impl HunkDecoder {
pub(crate) fn new() -> Self;
pub(crate) fn push(&mut self, chunk: Bytes);
pub(crate) fn buffered(&self) -> usize;
pub(crate) fn next(&mut self) -> io::Result<Option<Bytes>>;
pub(crate) fn next_multi(&mut self) -> io::Result<Option<Vec<Bytes>>>;
}
protocols/src/transports/grpc/stream.rs
pub struct GrpcStream {
send: SendStream<Bytes>,
recv: Recv,
mode: GrpcMode,
is_server: bool,
decoder: HunkDecoder,
ready: VecDeque<Bytes>,
write_pending: Option<Bytes>,
finished: bool,
_driver: Option<AbortOnDropHandle<()>>,
_guard: Option<StreamGuard>,
}
enum Recv {
Awaiting(ResponseFuture),
Body(RecvStream),
Done,
}
impl GrpcStream {
pub(crate) fn served(
send: SendStream<Bytes>,
recv: RecvStream,
mode: GrpcMode,
count: &Arc<StreamCount>,
) -> Self;
pub async fn connect<T>(
io: T,
secure: bool,
authority: &str,
service: &str,
mode: GrpcMode,
user_agent: Option<&str>,
) -> io::Result<Self>
where
T: AsyncRead + AsyncWrite + Unpin + Send + 'static;
}
pub(crate) struct StreamCount {
live: AtomicUsize,
changed: Notify,
}

A served GrpcStream carries a StreamGuard and no driver; a dialed one carries a driver and no guard.

GrpcPaths::new builds both paths from the service name, inserted verbatim with no normalisation. classify compares the request path exactly.

Path GrpcMode Message type
/<service>/Tun Gun one Hunk per gRPC message
/<service>/TunMulti Multi one MultiHunk per gRPC message
anything else none stream reset with REFUSED_STREAM
Message Headers
Client request (grpc_request) POST, URI https://<authority>/<service>/Tun (or http://, or /TunMulti), content-type: application/grpc, te: trailers, and user-agent when one is set
Server response (grpc_response) 200, content-type: application/grpc, sent as soon as the stream is accepted
Server close (finish_stream) trailers grpc-status: 0
Client close (finish_stream) an empty DATA frame with END_STREAM

DEFAULT_USER_AGENT is a desktop Chrome string, Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/152.0.7977.42 Safari/537.36. TransportKind::grpc sets it; TransportKind::user_agent(None) removes the header.

Every gRPC message on the stream is a length-prefixed message:

Offset Size Field Meaning
0 1 compressed flag Written as 0x00: uncompressed. The decoder does not read this byte
1 4 message length u32, big-endian: the length of the protobuf message that follows
5 message length message A protobuf Hunk or MultiHunk

A Hunk message is one length-delimited field:

Size Field Meaning
1 tag 0x0A (HUNK_DATA_TAG): field number 1, wire type 2
1 to 10 length Base-128 varint: the length of data
length data The tunnelled bytes

A MultiHunk message is zero or more of the same 0x0A / varint / data triples back to back. encode_multi_hunk skips empty payloads and computes the message length over the non-empty ones only.

Two encodings from the tests, byte for byte:

Call Bytes
encode_hunk(b"ok") 00 00 00 00 04 0A 02 6F 6B
encode_multi_hunk(["a", "", "bc"]) 00 00 00 00 07 0A 01 61 0A 02 62 63

A 130-byte payload needs a two-byte varint (82 01), so its message length is 0x85 (133).

GrpcStream::poll_write encodes each non-empty write as one message: one Hunk in Gun mode, a MultiHunk with a single entry in Multi mode. Encoding fails with hunk message too large only when the message length does not fit in a u32. The encoder does not apply MAX_GRPC_MESSAGE_LEN, so, as with WebSocket, a single write of more than about 1 MiB produces a message that an Etemenanki peer’s decoder rejects.

HTTP/2 DATA frames do not line up with gRPC messages, so HunkDecoder keeps the received chunks in a VecDeque<Bytes> with a running length and never copies until it has a whole message:

  1. take_frame peeks the 5-byte header across chunk boundaries (copy_prefix::<5>). Fewer than 5 bytes: Ok(None).
  2. It reads the length from bytes 1 to 4 and rejects anything over MAX_GRPC_MESSAGE_LEN immediately, before the body has arrived, with grpc message length exceeds the accepted maximum.
  3. With fewer than 5 + length bytes buffered it returns Ok(None) and waits.
  4. Otherwise it drops the header and splits off the body. When the body sits in one chunk this is a zero-copy split_to; otherwise the pieces are copied into one BytesMut.
  5. next (gun) requires the first byte to be 0x0A, reads the varint and splits off data. An empty message or empty data is skipped and the loop moves to the next message. Bytes after the declared data in the same message are discarded.
  6. next_multi walks every field of the message, requiring each to be 0x0A, and returns the non-empty entries. It can return an empty Vec.
Condition Error text (io::ErrorKind::InvalidData)
declared length over 1 MiB grpc message length exceeds the accepted maximum
length plus header overflows usize grpc frame too large
field tag other than 0x0A unexpected protobuf field in Hunk / unexpected protobuf field in MultiHunk
varint runs off the end truncated varint
varint that continues past its tenth byte varint overflow
varint does not fit usize hunk data too large
fewer bytes than the varint declares hunk data shorter than declared
Constant Value Client builder Server builder
H2_INITIAL_STREAM_WINDOW_SIZE 4 MiB yes yes
H2_INITIAL_CONNECTION_WINDOW_SIZE 16 MiB yes yes
H2_MAX_FRAME_SIZE 256 KiB yes yes
H2_MAX_CONCURRENT_STREAMS 256 no yes
MAX_GRPC_MESSAGE_LEN 1 MiB decoder decoder
H2_IDLE_TIMEOUT 300 s no Liveness
H2_KEEPALIVE_INTERVAL 60 s no Liveness
H2_KEEPALIVE_TIMEOUT 20 s no Liveness
DEFAULT_USER_AGENT Chrome string TransportKind::grpc no

The windows are sized for tunnel throughput over high-RTT paths. Everything else about the h2 builders is the h2 crate’s default.

sequenceDiagram
  participant CC as client core
  participant CG as GrpcStream client
  participant D as driver task
  participant S as serve_h2
  participant SG as GrpcStream served
  participant SC as server core
  CG->>D: handshake, then tokio::spawn(conn)
  CG->>S: HEADERS POST /svc/Tun
  S->>S: GrpcPaths::classify gives Gun
  S-->>CG: HEADERS 200 application/grpc
  S->>SG: GrpcStream::served, guard taken
  S->>SC: sink(TransportStream::Grpc)
  CC->>CG: poll_write payload
  CG->>SG: DATA encode_hunk(payload)
  SG->>SC: HunkDecoder::next, release_capacity
  SC->>SG: poll_write reply
  SG->>CG: DATA encode_hunk(reply)
  CG->>CC: Recv Awaiting then Body, decode
  CC->>CG: poll_shutdown
  CG->>SG: empty DATA with END_STREAM
  SC->>SG: poll_shutdown
  SG->>CG: trailers grpc-status 0

The client does not wait for the response headers in connect. Recv::Awaiting holds the ResponseFuture, and the first poll_read resolves it into Recv::Body, so a client-first protocol can write immediately after connect returns.

InboundTransport::accept for Grpc does TLS (optional) and the h2 server handshake inside the 10 s within("grpc", …) bound, then calls serve_h2:

protocols/src/transports/accept.rs
async fn serve_h2<T, F>(
mut conn: Connection<T, Bytes>,
paths: &GrpcPaths,
sink: &mut F,
) -> io::Result<()>
where
T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin,
F: FnMut(Accepted),

h2::server::Connection::accept both yields new streams and drives the connection’s I/O. Served streams live in other tasks, but their bytes only move while serve_h2 keeps polling conn.accept(). The loop is a tokio::select! over three arms:

flowchart TB
  top["loop: idle = count.is_idle()"] --> sel{"select!"}
  sel -->|"conn.accept() gives a stream"| cls{"classify path"}
  cls -->|"None"| rst["send_reset REFUSED_STREAM"]
  cls -->|"Gun or Multi"| resp["send_response 200"]
  resp --> sink["sink(GrpcStream::served)"]
  sel -->|"count.changed(), only if not idle"| prog["liveness.note_progress()"]
  sel -->|"liveness.watch(idle) resolves"| stop["break: drop the connection"]
  sel -->|"conn.accept() gives None"| stop
  rst --> top
  sink --> top
  prog --> top
  • Every accepted stream, refused or not, calls note_progress, restarting the idle deadline.
  • A stream whose path classify rejects is reset with REFUSED_STREAM; the connection keeps serving the others.
  • If send_response fails, the stream is dropped and the loop continues.
  • An error from conn.accept() ends serve_h2 with that error (io::Error::other).
  • StreamCount tracks live served streams. Each GrpcStream::served takes a StreamGuard that increments live; its Drop decrements it and calls notify_waiters, which wakes the count.changed() arm while that arm is waiting.

When serve_h2 returns, the Connection is dropped and every stream still open on it fails on its next read or write.

Liveness supervises a served connection against two hazards, each with its own deadline:

protocols/src/transports/grpc/liveness.rs
pub(crate) struct Liveness {
ping_pong: Option<PingPong>,
next_ping: Instant,
pong_due: Option<Instant>,
idle_due: Instant,
}
impl Liveness {
pub(crate) fn new(ping_pong: Option<PingPong>) -> Self;
pub(crate) fn note_progress(&mut self);
pub(crate) async fn watch(&mut self, idle: bool);
}
pub(crate) enum Verdict {
Alive,
Dead,
}
Hazard Deadline Applies when Reset by
A peer completes the handshake and opens no stream, which would park Connection::accept forever idle_due, H2_IDLE_TIMEOUT = 300 s idle is true: StreamCount is zero note_progress: a stream accepted or a stream ended
A peer vanishes without a FIN while streams are open, so the idle deadline never applies pong_due, H2_KEEPALIVE_TIMEOUT = 20 s after each PING a PING is outstanding the matching PONG

next_ping fires every H2_KEEPALIVE_INTERVAL (60 s) from the moment serve_h2 starts, independent of traffic. When it fires and no PING is outstanding, tick sends Ping::opaque() and sets pong_due in the same step. A failed send_ping is Dead at once. poll_pong resolving Ok clears pong_due; resolving Err is Dead. serve_h2 passes conn.ping_pong(), which h2 hands out once; with None, poll_pong pends forever and only the timers apply.

watch loops over tick until the verdict is Dead rather than resolving on every tick. In tokio::select!, an arm whose pattern does not match is disabled for the rest of that call, so an arm that resolved on every tick would stop re-arming its timers the first time it found the peer alive. The loop is cancel-safe, which it must be because every other arm of serve_h2’s select! drops it: its awaits are a timer and a poll, no state changes across an await, and pong_due is set in the same synchronous step that sends the PING, so a cancelled watch never leaves a PING unwatched.

A Dead verdict ends serve_h2 with Ok(()): giving up on a peer is not an error.

GrpcStream::connect opens a fresh HTTP/2 connection for every dial and one stream on it:

  1. configured_client_builder().handshake(io) returns (SendRequest, Connection).
  2. The Connection future is spawned with tokio::spawn and wrapped in AbortOnDropHandle. When it ends with an error, the task logs grpc connection ended: … at debug level.
  3. send_req.ready() waits for stream capacity, and send_request(request, false) opens the stream with the body still open.
  4. The SendRequest handle is dropped when connect returns; the open stream keeps the connection alive.

The driver handle is stored in _driver. Dropping the GrpcStream drops the handle, which aborts the driver task and with it the connection. A dialed connection is never shared between streams, and it runs no Liveness.

write_pending: Option<Bytes> holds at most one encoded message:

  1. poll_write first drains the previous message with poll_drain. It returns Pending until that is fully handed to h2, which is how HTTP/2 flow control pushes back on the writer.
  2. It encodes the new buffer, stores it in write_pending and makes one best-effort poll_drain.
  3. It returns Ok(buf.len()) even if part of the message is still pending. poll_flush and the next poll_write finish it.

poll_drain calls reserve_capacity(pending.len()), waits on poll_capacity, and sends as many bytes as the granted capacity allows with send_data(chunk, false). A None from poll_capacity means the stream is closed: BrokenPipe, h2 send stream closed. poll_shutdown drains, then calls finish_stream once (guarded by finished).

poll_read works in a fixed order:

  1. Hand out bytes from the front of ready if there are any.
  2. Otherwise run decode, which asks the decoder for one message (or one MultiHunk batch) and appends its payloads to ready.
  3. Otherwise pull one DATA chunk from h2 (poll_data) into the decoder, or resolve Recv::Awaiting first. None from poll_data moves to Recv::Done, which is EOF.

The decoder is fed one DATA chunk at a time, and only when ready is empty and no complete message is buffered, so nothing is pulled from h2 until the reader asks for it.

decode returns window credit with flow_control().release_capacity(consumed), where consumed is the drop in HunkDecoder::buffered() across the call. Bytes of an unfinished message keep holding their window credit, so a peer cannot send more of one message than the stream window allows before the reader consumes it, on top of the MAX_GRPC_MESSAGE_LEN cap on the declared length.

Invariant Enforced by Pinned by
One non-empty write is exactly one WebSocket Binary message WsStream::poll_write sends Message::Binary(Bytes::copy_from_slice(buf)) no dedicated test; the pipeline round trips only compare the echoed bytes
One non-empty write is exactly one gRPC message GrpcStream::poll_write → encode_hunk / encode_multi_hunk encode_hunk_writes_exact_small_frame, encode_multi_hunk_writes_exact_repeated_fields (protocols/tests/unit/transports/grpc_framing.rs) pin the encoding; one_grpc_connection_carries_many_streams (protocols/tests/pipeline/transports.rs) checks that a served stream echoes a 10-byte read as exactly one Hunk
The same path string works on both sides; ?ed= never reaches the wire parse_early_data_path + normalize_path in WsRoute::new and WsTarget::new normalizes_paths, parses_outbound_early_data_path (protocols/tests/unit/transports/ws_endpoint.rs)
At most min(ed, MAX_EARLY_DATA) bytes ride the upgrade, and they are the first bytes the server core reads start_upgrade returns take; from_upgraded seeds pending_read ws_early_data_is_the_first_bytes_the_server_reads (protocols/tests/pipeline/transports.rs)
Early data over 16 KiB is refused before the upgrade decode_early_data_header returns Err(()), callback answers 413 rejects_oversized_early_data_header, callback_rejects_oversized_early_data (ws_endpoint.rs)
Valid early data is echoed; invalid early data is ignored, not echoed WsCallback::on_request callback_captures_and_echoes_valid_early_data, callback_ignores_invalid_early_data, decodes_xray_early_data_header, encodes_xray_early_data_header (ws_endpoint.rs)
A wrong path or Host never upgrades WsRoute::accepts → reject_404 route_accepts_only_matching_path_and_host, host_matching_ignores_case_and_port (ws_endpoint.rs), ws_rejects_a_wrong_path_at_the_upgrade (transports.rs)
No received WebSocket message or frame over 1 MiB is accepted, on inbound and dialed sessions alike ws_config passed to both tungstenite entry points the_configured_limits_replace_tungstenite_defaults (ws_endpoint.rs) checks the config values
A declared gRPC length over 1 MiB is refused before its body is buffered HunkDecoder::take_frame checks MAX_GRPC_MESSAGE_LEN on the header decoder_rejects_a_message_longer_than_the_maximum, decoder_rejects_an_oversized_length_before_buffering_the_body, decoder_accepts_a_message_at_the_maximum (grpc_framing.rs)
DATA frame boundaries do not matter HunkDecoder chunk queue and copy_prefix decoder_waits_for_split_frame, decoder_reads_multiple_frames_from_one_chunk, decoder_reads_multi_hunk_entries (grpc_framing.rs)
Malformed hunks fail with InvalidData tag and length checks in next / next_multi decoder_rejects_unexpected_field, decoder_rejects_short_declared_payload (grpc_framing.rs)
Only /<service>/Tun and /<service>/TunMulti are served GrpcPaths::classify, REFUSED_STREAM otherwise classify_path_matches_tun_modes (protocols/tests/unit/transports/grpc_liveness.rs) pins the classification; no test drives the reset
One served connection carries many streams serve_h2 loop, sink per stream one_grpc_connection_carries_many_streams (transports.rs)
A served connection with no streams is dropped after H2_IDLE_TIMEOUT Liveness::idle_due liveness_gives_up_on_a_connection_with_no_streams, liveness_restarts_the_idle_deadline_on_progress (grpc_liveness.rs), a_connection_that_opens_no_stream_is_given_up_on (protocols/tests/unit/transports/accept.rs)
A served connection with live streams is not dropped by the idle deadline idle_due is only consulted when idle is true liveness_keeps_a_connection_carrying_streams (grpc_liveness.rs)
Window credit is returned only for consumed bytes GrpcStream::decode releases held - buffered() no dedicated test
A dialed connection dies with its stream AbortOnDropHandle in _driver no dedicated test
Liveness::watch is cancel-safe pong_due set in the same step as send_ping no dedicated test
Situation What happens
Inbound TLS plus WebSocket upgrade, or TLS plus h2 handshake, exceeds TRANSPORT_HANDSHAKE_TIMEOUT accept fails with TimedOut, websocket handshake timed out or grpc handshake timed out; the app logs inbound transport failed at debug level
Upgrade refused (404, 413) The server’s accept fails. A client dialed without ed fails in connect with an error whose text contains the status; a deferred client stores it in State::Failed and reports it on its next call
Deferred upgrade fails State::Failed; the first write already returned Ok(take), every later call returns the stored error
WebSocket peer closes or disappears EOF on read (Close, stream end, reset without closing handshake)
WebSocket quiet for 300 s EOF on read
Write after the WebSocket is closed BrokenPipe, websocket closed
Unknown gRPC path That stream is reset with REFUSED_STREAM; the connection continues
send_response fails That stream is dropped silently; the connection continues
conn.accept() errors serve_h2 returns the error and the connection is dropped
Liveness judges the peer dead serve_h2 returns Ok(()), the connection is dropped, open streams fail
h2 stream or connection error on a GrpcStream h2_err: an I/O error is unwrapped, anything else becomes io::Error::other
Send side closed while draining BrokenPipe, h2 send stream closed
Malformed gRPC message InvalidData from the decoder, surfaced from poll_read
Dialed GrpcStream dropped The driver task is aborted and the HTTP/2 connection closes
Dialed WsStream dropped while upgrading The boxed handshake future and its socket are dropped
Constant Value Where
TRANSPORT_HANDSHAKE_TIMEOUT 10 s accept.rs: TLS plus upgrade or h2 handshake on accept
MAX_EARLY_DATA 16 KiB ws/endpoint.rs: decoded early data, and the cap on a client’s ed
MAX_WS_MESSAGE_LEN 1 MiB ws/endpoint.rs: tungstenite max_message_size and max_frame_size
WS_KEEPALIVE_INTERVAL 60 s ws/stream.rs: quiet period before a Ping
WS_IDLE_TIMEOUT 300 s ws/stream.rs: quiet period before EOF
MAX_GRPC_MESSAGE_LEN 1 MiB grpc/settings.rs: largest declared gRPC message
H2_INITIAL_STREAM_WINDOW_SIZE 4 MiB grpc/settings.rs
H2_INITIAL_CONNECTION_WINDOW_SIZE 16 MiB grpc/settings.rs
H2_MAX_FRAME_SIZE 256 KiB grpc/settings.rs
H2_MAX_CONCURRENT_STREAMS 256 grpc/settings.rs: server only
H2_IDLE_TIMEOUT 300 s grpc/settings.rs: served connection with no streams
H2_KEEPALIVE_INTERVAL 60 s grpc/settings.rs: PING period on a served connection
H2_KEEPALIVE_TIMEOUT 20 s grpc/settings.rs: PONG deadline

TCP keepalive (TCP_KEEPALIVE_IDLE 120 s, TCP_KEEPALIVE_INTERVAL 30 s, TCP_KEEPALIVE_RETRIES 3) sits under both carriers; see Transports: TCP and TLS.

File Tests Covers
protocols/tests/unit/transports/ws_endpoint.rs normalizes_paths, parses_outbound_early_data_path, host_matching_ignores_case_and_port, route_accepts_only_matching_path_and_host, decodes_xray_early_data_header, encodes_xray_early_data_header, rejects_oversized_early_data_header, callback_captures_and_echoes_valid_early_data, callback_ignores_invalid_early_data, callback_rejects_oversized_early_data, the_configured_limits_replace_tungstenite_defaults Path and ?ed= handling, Host check, early-data codec and callback, size limits
protocols/tests/unit/transports/grpc_framing.rs encode_hunk_writes_exact_small_frame, encode_hunk_writes_multi_byte_varint_length, encode_multi_hunk_writes_exact_repeated_fields, decoder_waits_for_split_frame, decoder_reads_multiple_frames_from_one_chunk, decoder_reads_multi_hunk_entries, decoder_rejects_unexpected_field, decoder_rejects_short_declared_payload, decoder_rejects_a_message_longer_than_the_maximum, decoder_rejects_an_oversized_length_before_buffering_the_body, decoder_accepts_a_message_at_the_maximum Exact wire bytes, reassembly, malformed input, the 1 MiB cap
protocols/tests/unit/transports/grpc_liveness.rs classify_path_matches_tun_modes, liveness_gives_up_on_a_connection_with_no_streams, liveness_keeps_a_connection_carrying_streams, liveness_restarts_the_idle_deadline_on_progress Path classification; the idle deadline on a paused clock
protocols/tests/unit/transports/accept.rs a_connection_that_opens_no_stream_is_given_up_on serve_h2 and Liveness wired together over tokio::io::duplex, with the clock stepped by hand
protocols/tests/pipeline/transports.rs ws_round_trip, ws_over_tls_with_early_data_round_trip, ws_early_data_is_the_first_bytes_the_server_reads, ws_rejects_a_wrong_path_at_the_upgrade, grpc_round_trip, grpc_multi_mode_over_tls_round_trip, one_grpc_connection_carries_many_streams The four _round_trip tests echo hello and 200 000 bytes through InboundTransport and TransportConnector, then expect a clean EOF. The early-data test writes 19 bytes with ?ed=16 and checks that the first write returns 16. The wrong-path test expects 404 in the dial error. The many-streams test opens three streams on one raw h2 client connection
app/tests/integration/e2e_xray.rs app_client_ws_xray_server_plain, app_server_ws_xray_client_plain, app_client_grpc_xray_server_plain, app_server_grpc_xray_client_plain, app_client_ws_xray_server_tls, app_client_grpc_xray_server_tls, app_server_ws_xray_client_tls, app_server_grpc_xray_client_tls VLESS over WebSocket and gRPC against a real Xray binary, both directions
app/tests/integration/e2e_xray_vmess.rs app_server_vmess_grpc_xray_client_tls, app_client_vmess_ws_xray_server_early_data_plain, app_server_vmess_ws_xray_client_early_data_plain gRPC with TLS and WebSocket early data (/vmess?ed=2048) against Xray
app/tests/integration/e2e_xray_mux.rs vless_mux_over_ws_tls mux.cool composed over WebSocket with TLS

The Xray interop tests build Xray with go and skip themselves, printing SKIP, when go is not available or the build fails. No test drives WS_KEEPALIVE_INTERVAL or WS_IDLE_TIMEOUT, or the PING/PONG half of Liveness; a change there needs a paused-clock test in the style of grpc_liveness.rs. See Testing for how the suites are run.