Transports: WebSocket and gRPC
Source files: 22 · checked against Etemenanki 596916d · katana v3.0.1
Etemenanki/protocols/src/transports/ws/endpoint.rsEtemenanki/protocols/src/transports/ws/stream.rsEtemenanki/protocols/src/transports/grpc/framing.rsEtemenanki/protocols/src/transports/grpc/liveness.rsEtemenanki/protocols/src/transports/grpc/settings.rsEtemenanki/protocols/src/transports/grpc/stream.rsEtemenanki/protocols/src/transports/accept.rsEtemenanki/protocols/src/transports/connect.rsEtemenanki/protocols/src/transports/stream.rsEtemenanki/app/src/transport.rsEtemenanki/app/src/inbound/mod.rsEtemenanki/app/src/outbound/mod.rsEtemenanki/app/src/serve.rsEtemenanki/protocols/tests/unit/transports/ws_endpoint.rsEtemenanki/protocols/tests/unit/transports/grpc_framing.rsEtemenanki/protocols/tests/unit/transports/grpc_liveness.rsEtemenanki/protocols/tests/unit/transports/accept.rsEtemenanki/protocols/tests/pipeline/transports.rsEtemenanki/app/tests/integration/e2e_xray.rsEtemenanki/app/tests/integration/e2e_xray_vmess.rsEtemenanki/app/tests/integration/e2e_xray_mux.rskatana/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.
Responsibilities
Section titled “Responsibilities”| 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.
Entry points
Section titled “Entry points”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.
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);}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.
How the programs build them
Section titled “How the programs build them”| 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.
WebSocket
Section titled “WebSocket”Key types
Section titled “Key types”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.
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.
Path normalisation and ?ed=
Section titled “Path normalisation and ?ed=”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
edwith 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 atNone. - When
edis present, empty query parts are dropped and the remaining ones are re-joined with&. Withouted, the path is returned untouched. WsStream::connectonly defers the upgrade when the limit is greater than zero, so?ed=0means no early data.
Host validation
Section titled “Host validation”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.
Early data
Section titled “Early data”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.
WsStream states
Section titled “WsStream states”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.
Reading and writing
Section titled “Reading and writing”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.
Keepalive and idle timeout
Section titled “Keepalive and idle timeout”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.
Message and frame limits
Section titled “Message and frame limits”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.
Key types
Section titled “Key types”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);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>>>;}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.
Paths and headers
Section titled “Paths and headers”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.
Wire format
Section titled “Wire format”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.
Decoding
Section titled “Decoding”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:
take_framepeeks the 5-byte header across chunk boundaries (copy_prefix::<5>). Fewer than 5 bytes:Ok(None).- It reads the length from bytes 1 to 4 and rejects anything over
MAX_GRPC_MESSAGE_LENimmediately, before the body has arrived, withgrpc message length exceeds the accepted maximum. - With fewer than
5 + lengthbytes buffered it returnsOk(None)and waits. - 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 oneBytesMut. next(gun) requires the first byte to be0x0A, reads the varint and splits offdata. An empty message or emptydatais skipped and the loop moves to the next message. Bytes after the declareddatain the same message are discarded.next_multiwalks every field of the message, requiring each to be0x0A, and returns the non-empty entries. It can return an emptyVec.
| 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 |
HTTP/2 settings
Section titled “HTTP/2 settings”| 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.
One HTTP/2 stream, end to end
Section titled “One HTTP/2 stream, end to end”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.
Serving an HTTP/2 connection
Section titled “Serving an HTTP/2 connection”InboundTransport::accept for Grpc does TLS (optional) and the h2 server handshake inside the 10 s within("grpc", …) bound, then calls serve_h2:
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
classifyrejects is reset withREFUSED_STREAM; the connection keeps serving the others. - If
send_responsefails, the stream is dropped and the loop continues. - An error from
conn.accept()endsserve_h2with that error (io::Error::other). StreamCounttracks live served streams. EachGrpcStream::servedtakes aStreamGuardthat incrementslive; itsDropdecrements it and callsnotify_waiters, which wakes thecount.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
Section titled “Liveness”Liveness supervises a served connection against two hazards, each with its own deadline:
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.
The client stream owns its connection
Section titled “The client stream owns its connection”GrpcStream::connect opens a fresh HTTP/2 connection for every dial and one stream on it:
configured_client_builder().handshake(io)returns(SendRequest, Connection).- The
Connectionfuture is spawned withtokio::spawnand wrapped inAbortOnDropHandle. When it ends with an error, the task logsgrpc connection ended: …at debug level. send_req.ready()waits for stream capacity, andsend_request(request, false)opens the stream with the body still open.- The
SendRequesthandle is dropped whenconnectreturns; 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.
Writes under flow control
Section titled “Writes under flow control”write_pending: Option<Bytes> holds at most one encoded message:
poll_writefirst drains the previous message withpoll_drain. It returnsPendinguntil that is fully handed to h2, which is how HTTP/2 flow control pushes back on the writer.- It encodes the new buffer, stores it in
write_pendingand makes one best-effortpoll_drain. - It returns
Ok(buf.len())even if part of the message is still pending.poll_flushand the nextpoll_writefinish 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).
Reads and window credit
Section titled “Reads and window credit”poll_read works in a fixed order:
- Hand out bytes from the front of
readyif there are any. - Otherwise run
decode, which asks the decoder for one message (or oneMultiHunkbatch) and appends its payloads toready. - Otherwise pull one DATA chunk from h2 (
poll_data) into the decoder, or resolveRecv::Awaitingfirst.Nonefrompoll_datamoves toRecv::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.
Invariants
Section titled “Invariants”| 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 |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| 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 |
Limits
Section titled “Limits”| 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.