Skip to content

Links, connectors and net types

Source files: 50 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/net.rs
  • Etemenanki/concepts/src/sniff.rs
  • Etemenanki/concepts/src/relay.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/tests/runtime.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/environment/src/dial/mod.rs
  • Etemenanki/environment/src/dial/tcp.rs
  • Etemenanki/environment/src/dial/udp.rs
  • Etemenanki/environment/src/routing.rs
  • Etemenanki/environment/tests/integration/udp.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/socks/handshake.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/app/src/connector.rs
  • Etemenanki/app/src/router.rs
  • Etemenanki/app/src/outbound/freedom.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • Etemenanki/app/src/outbound/udp_fanout.rs
  • Etemenanki/app/src/flow.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/src/hysteria/server/authenticator.rs
  • Etemenanki/protocols/tests/pipeline/transports.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/protocols/tests/unit/socks/server.rs
  • Etemenanki/protocols/tests/unit/socks/protocol.rs
  • Etemenanki/app/tests/integration/e2e_sniff.rs
  • katana/src/connector.rs
  • katana/src/router.rs
  • katana/src/traffic.rs
  • katana/src/outbound/proxy.rs
  • katana/tests/unit/serve.rs
  • katana/tests/unit/e2e.rs
  • katana/tests/unit/runtime.rs

Every crate above etemenanki-concepts speaks one small, shared vocabulary. link.rs says what an outbound is and how one is dialed. net.rs says what an address and a user are. sniff.rs says what sniffing found. relay.rs holds two byte-copy futures for paths with no protocol in between. None of these modules resolves names, applies socket policy or routes. They fix the shapes that the server runtime, the client runtime, the dialers, the protocol cores, the app and katana plug into.

This page is for contributors who add a protocol, an outbound kind or a connector, or who change how the runtime treats links. It gives each type with its exact signature, says who implements it and who consumes it, and ties every rule to the mechanism that enforces it and the test that pins it.

Module Defines Leaves to others
concepts/src/link.rs DatagramLink, UdpOutbound, Outbound<S, D>, Connector, SocketConnector, SocketTarget, NoStream, NoDatagram Socket options, connect timeouts and DNS. The dialers in etemenanki-environment and the application outbounds own those.
concepts/src/net.rs UserAuthorization, NetworkUser<T>, DialNetwork, Remote, Destination Resolution. A Destination may name a domain, and whoever reaches it decides how to resolve it.
concepts/src/sniff.rs SniffedProtocol, SniffedBehavior, Sniffer The sniffers (TlsSniffer, HttpSniffer), their byte budget and deadline. They live in protocols/src/sniff/.
concepts/src/relay.rs UnidirectionalConnection, BidirectionalConnection, Relayed Timeouts and accounting. A caller wraps the future.

Flow<T> in protocols/src/flow.rs is not part of the concepts crate. It appears on this page because it is the Target every server core in etemenanki-protocols hands its connector, and so it is where the net and sniff types meet:

protocols/src/flow.rs
pub struct Flow<T> {
pub destination: Destination,
pub user: NetworkUser<T>,
pub sniffed: Option<SniffedBehavior>,
pub source: Option<IpAddr>,
}
impl<T> Flow<T> {
pub fn new(destination: Destination, user: NetworkUser<T>, source: Option<IpAddr>) -> Self;
pub fn toward(&self, destination: Destination) -> Self;
}
impl<T> Clone for Flow<T>;

A server core decides what to open. The runtime asks a Connector to open it and gets back an Outbound: either a byte stream or a DatagramLink. From then on the runtime moves bytes between the core and that link.

flowchart LR
  wire["client transport"] --> rt["ProxyServerRuntime"]
  rt -- "Event" --> core["ProxyCoreDecode"]
  core -- "Effect::Open, target = Flow" --> rt
  sn["Sniffer result: SniffedBehavior"] -.-> |"Flow.sniffed"| core
  rt -- "Connector::connect(target)" --> conn["Connector"]
  conn -- "Outbound::Stream" --> s["AsyncRead + AsyncWrite + Unpin"]
  conn -- "Outbound::Datagram" --> d["DatagramLink, Addr = Destination"]
  s --> dest["destination or upstream"]
  d --> dest

The same traits recur one level down. ProxyClientConnector is a Connector whose outbounds are ProxyClientRuntimes, and each of those dials its upstream through another Connector: TransportConnector in both the app and katana. A connection routed to a proxy outbound therefore runs three connectors inside one task: the routing connector (AppConnector or KatanaConnector), the outbound’s ProxyClientConnector and its TransportConnector.

Layer Implementor Target Stream / Datagram
App routing app/src/connector.rs → AppConnector Flow (the app’s Flow<()>) OutboundStream / FanOutLink
katana routing src/connector.rs → KatanaConnector Flow<UserTag> Metered<OutboundStream> / FanOut
Direct outbound app/src/outbound/freedom.rs → FreedomConnector Flow TcpStream / ResolvingUdp
Proxy outbound concepts/src/client.rs → ProxyClientConnector whatever its Make closure takes: the app’s Flow, or a Destination in katana ProxyClientRuntime<BUF_SIZE, S, Conn> / ProxyClientRuntime<BUF_SIZE, D, Conn>
Upstream dial protocols/src/transports/connect.rs → TransportConnector Destination TransportStream / NoDatagram
Tunnel outbounds protocols/src/wireguard/connector.rs → WgConnector, protocols/src/hysteria/connector.rs → Hy2Connector Flow<T> WgStream / WgDatagramLink, Hy2Stream / Hy2DatagramLink
Host sockets environment/src/dial/mod.rs → Dialer DialTarget or SocketTarget TcpStream / DualStackUdp or UdpOutbound
Reference and tests concepts/src/link.rs → SocketConnector, and any closure anything anything
concepts/src/link.rs
pub trait DatagramLink: Unpin {
type Addr;
fn poll_send_to(
&mut self,
cx: &mut Context<'_>,
buf: &[u8],
to: &Self::Addr,
) -> Poll<io::Result<usize>>;
fn poll_recv_from(
&mut self,
cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<io::Result<Self::Addr>>;
}

A datagram link is packet-oriented. One poll_send_to is one datagram to to, and one poll_recv_from is one datagram from the address it returns. There is no end of stream, no flush and no shutdown. The methods take &mut self rather than Pin<&mut Self>, and the trait requires Unpin:

  • The server runtime stores outbounds by value in a BTreeMap and moves them when the map rebalances. Unpin makes that sound, and the module docs of link.rs give it as the reason for the bound.
  • The trait stays object safe once Addr is fixed. app/src/outbound/proxy.rs → OutboundDatagram wraps a Box<dyn DatagramLink<Addr = Destination> + Send>, and katana’s src/outbound/proxy.rs has a type of the same name.

Addr is what a packet is addressed by, and it decides which of two roles a link can play:

Role Required Addr Enforced by Implementors
Outbound, behind a Connector Destination Conn::Datagram: DatagramLink<Addr = Destination> on the ProxyServerRuntime impls that poll it, including its Future and Stream impls UdpOutbound, DualStackUdp, ResolvingUdp, FanOutLink, OutboundDatagram, BlackholeLink, SocksUdpLink<S>, WgDatagramLink, Hy2DatagramLink, ProxyClientRuntime over a datagram codec, katana’s FanOut
Transport, the client side of a runtime built with over_datagrams the core’s TransportAddr D: DatagramLink<Addr = Core::TransportAddr> on ProxyServerRuntime::over_datagrams TunUdpLink (SocketAddr, the TUN UDP path), QuicDatagrams ((), the datagrams of one Hysteria 2 QUIC connection), UdpSocket (SocketAddr, used by the concepts tests)

The two bounds, as the runtime writes them:

concepts/src/runtime.rs
impl<const BUF_SIZE: usize, Core, D, Conn>
ProxyServerRuntime<BUF_SIZE, Core, DatagramTransport<D>, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
D: DatagramLink<Addr = Core::TransportAddr>,
Conn: Connector<Core::Target>,
{
pub fn over_datagrams(link: D, core: Core, connector: Conn) -> Self;
}
impl<const BUF_SIZE: usize, Core, Trans, Conn> Future
for ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>
where
Core: ProxyCoreDecode,
Trans: Transport<Addr = Core::TransportAddr>,
Conn: Connector<Core::Target>,
Conn::Datagram: DatagramLink<Addr = Destination>,
{
type Output = Result<Traffic, RuntimeError<Core::Error>>;
}

The runtime reads a link’s results differently in the two roles. A link author needs to know which result keeps the link alive:

Call and result As an outbound As a transport
poll_send_to → Pending The SendTo effect stays at the head of the queue and holds every effect behind it. The key’s waker resumes it. The packet stays staged.
poll_send_to → Ready(Ok(n)) The packet is done and n counts toward Traffic::outbound_tx. A short count is not retried. The packet is done.
poll_send_to → Ready(Err(e)) The packet is dropped and the key stays live. The core receives Event::SendFailed { key, to, error }. The packet is dropped and the core receives Event::TransportSendFailed { to, error }. The runtime keeps running.
poll_recv_from → Ready(Ok(from)) Event::Datagram { key, from, data }, one whole packet. Event::TransportDatagram { from, data }, one whole packet.
poll_recv_from → Ready(Err(e)) The key is removed and the core receives Event::OutboundError { key, error }. Fatal: the runtime ends with RuntimeError::Transport.

The runtime reads a packet from a datagram outbound into at most datagram_limit() bytes: the core’s MAX_DATAGRAM (default 4096), capped at BUF_SIZE - STAGING_RESERVE. From a datagram transport it reads at most MAX_DATAGRAM capped at BUF_SIZE. A longer packet is truncated in both roles, as a kernel recv truncates. The runtime polls a datagram outbound only when the staging buffer has STAGING_RESERVE + datagram_limit() bytes free, so a client that stops reading stalls the downlink at the socket rather than losing packet tails. The server runtime page covers the scheduling.

A link that drops a packet on purpose returns Ready(Ok(buf.len())). BlackholeLink does that for every packet, and ResolvingUdp does it for a name that does not resolve. A link that returns an error instead lets the core decide what to do.

concepts/src/link.rs
impl DatagramLink for UdpSocket {
type Addr = SocketAddr;
// poll_send_to -> UdpSocket::poll_send_to(self, cx, buf, *to)
// poll_recv_from -> UdpSocket::poll_recv_from(self, cx, buf)
}

The bare tokio socket forwards to its inherent methods. Its Addr is SocketAddr, so it fits only the transport role: one socket serving many peers, each named by a concrete address. It cannot be a connector’s Datagram type under the server runtime, because the runtime requires Addr = Destination there.

concepts/src/link.rs
pub struct UdpOutbound(pub UdpSocket);
impl DatagramLink for UdpOutbound {
type Addr = Destination;
}

UdpOutbound is the same socket in the outbound role. poll_send_to calls Destination::socket_addr():

  • For Remote::IpAddr, it sends to that address.
  • For Remote::Domain, it returns ErrorKind::Unsupported with the text a plain UDP outbound cannot resolve a domain. The runtime turns that into Event::SendFailed, and the key stays live.

poll_recv_from wraps the peer’s SocketAddr in Destination::udp, so Event::Datagram’s from is always an IP for this link.

The refusal is by design. Resolution is policy (which resolver answers, which address family a destination may use), and that policy belongs to the application. environment/src/dial/udp.rs → DualStackUdp makes the same choice with the text udp: a plain dual-stack link cannot resolve a domain. The link that resolves is app/src/outbound/freedom.rs → ResolvingUdp, which looks names up through the app’s Resolver; the app outbounds page describes it.

concepts/src/link.rs
pub enum Outbound<S, D> {
Stream(S),
Datagram(D),
}

Outbound<S, D> is a dialed outbound. The connector picks the variant, not the caller. AppConnector and KatanaConnector return Datagram for a flow whose destination.network is DialNetwork::Udp, without dialing anything, and Stream otherwise. (KatanaConnector first admits the flow’s user and fails the dial for a user that is no longer registered.) Consumers that expect one kind reject the other explicitly:

Consumer Receives Result
ProxyServerRuntime, applying Effect::Forward or Effect::ForwardHeld a datagram outbound RuntimeError::WrongLinkKind (stream effect on a datagram outbound or vice versa), which ends the connection
ProxyServerRuntime, applying Effect::SendTo or Effect::SendToHeld a stream outbound RuntimeError::WrongLinkKind
ProxyServerRuntime, applying Effect::Shutdown a datagram outbound not an error: the effect completes at once without calling the link, and the link stays open
ProxyClientRuntime, dialing the upstream wire Outbound::Datagram ErrorKind::Unsupported (proxy client runtime needs a stream to the upstream, got a datagram socket) and the wire goes Down
protocols/src/socks/server.rs, a CONNECT Outbound::Datagram ErrorKind::Unsupported (socks: a CONNECT was answered with a datagram link)
protocols/src/socks/server.rs, a UDP ASSOCIATE Outbound::Stream ErrorKind::Unsupported (socks: an association was answered with a stream)
proxy_stream in app/src/outbound/proxy.rs and in katana’s src/outbound/proxy.rs Outbound::Datagram ErrorKind::Other (a TCP flow was dialed as datagrams)
proxy_datagram, in the same two files Outbound::Stream ErrorKind::Other (a UDP flow was dialed as a stream)

The enum is also used as a plain “one of two kinds” type. ProxyClientConnector’s Make closure returns Outbound<S, D> where S and D are codecs, not links: it picks the stream codec or the datagram codec for a flow. The app and katana both have their own type named outbound::Outbound, so they import the concepts module and write this one as link::Outbound.

concepts/src/link.rs
pub trait Connector<Target> {
type Stream: AsyncRead + AsyncWrite + Unpin;
type Datagram: DatagramLink;
type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;
}

A Connector turns a Target into an Outbound. Under a server runtime, Target is the core’s ProxyCoreDecode::Target, which is Flow<T> for every core in etemenanki-protocols. Under a client runtime, it is the codec’s ProxyCoreEncodeHandshake::Target: the upstream server, a Destination in the app.

Properties of the trait that shape the code around it:

  • connect takes &mut self, so a connector may keep state between dials. The server runtime calls it synchronously while it applies an Effect::Open. ProxyClientRuntime::new borrows the connector (connector: &mut Conn) only for the call that creates the dial future.
  • Future has no Unpin or Send bound. Both runtimes box it once per open: LinkState::Connecting(Pin<Box<F>>) in the server runtime and Wire::Connecting(Pin<Box<F>>) in the client runtime. A connector may therefore return an unnamed async type or a std::future::Ready. Every production connector except ProxyClientConnector returns Pin<Box<dyn Future<Output = …> + Send>>, so that the runtime holding it can be spawned. ProxyClientConnector returns the named future ProxyClientConnecting, which resolves only once the upstream is dialed and the codec’s handshake is done.
  • The future’s error is the outbound’s failure. The server runtime removes the key, drops every effect still queued for it (forget_key), and delivers Event::ConnectFailed { key, error }.
concepts/src/link.rs
impl<Target, F, Fut, S, D> Connector<Target> for F
where
F: FnMut(Target) -> Fut,
Fut: Future<Output = io::Result<Outbound<S, D>>>,
S: AsyncRead + AsyncWrite + Unpin,
D: DatagramLink,
{
type Stream = S;
type Datagram = D;
type Future = Fut;
}

Any FnMut(Target) -> Fut whose future yields io::Result<Outbound<S, D>> is a connector, so a test never needs to name a type. That includes plain functions (a fn item is FnMut) and, through the standard library’s FnMut impl for &mut F, a mutable borrow of a closure. The concepts tests use closures and plain functions:

concepts/tests/runtime.rs
fn udp_connector(_: ()) -> Ready<io::Result<Outbound<NoStream, UdpOutbound>>>
fn dialer(dials: Dials) -> impl FnMut(String) -> DialFuture

dial_failure_surfaces_on_first_use in concepts/tests/client.rs uses a closure dial as the connector and passes &mut dial because ProxyClientRuntime::new takes connector: &mut Conn. katana’s tests build stream-only connectors the same way: a closure in tests/unit/serve.rs and named connector types in tests/unit/e2e.rs, both with Datagram = NoDatagram. One of the named types, TcpConnector (a Connector<Destination>), is also the connector of the VMess and VLESS test clients (vmess_client, vless_client) that tests/unit/runtime.rs reuses.

concepts/src/link.rs
pub struct SocketConnector;
pub enum SocketTarget {
Tcp(SocketAddr),
Udp {
ipv6: bool,
},
}
impl Connector<SocketTarget> for SocketConnector {
type Stream = TcpStream;
type Datagram = UdpOutbound;
type Future =
Pin<Box<dyn Future<Output = io::Result<Outbound<TcpStream, UdpOutbound>>> + Send>>;
}

SocketConnector is the dependency-free reference connector:

  • SocketTarget::Tcp(addr) runs TcpStream::connect(addr).
  • SocketTarget::Udp { ipv6 } binds an unconnected socket on [::]:0 or 0.0.0.0:0 and wraps it in UdpOutbound. The socket is never connected, so Effect::SendTo picks the peer per packet. Only the address family is fixed at bind time.

Nothing in the workspace or in katana constructs SocketConnector at the pinned revisions. The policy-aware implementation of the same contract is environment/src/dial/mod.rs → Dialer, which implements Connector<SocketTarget> with the same Stream and Datagram types:

SocketConnector Dialer as Connector<SocketTarget>
TCP TcpStream::connect(addr), no timeout of its own TcpDialer::connect(addr), with SocketOptions and the connect timeout (DEFAULT_CONNECT_TIMEOUT, 10 s)
UDP UdpSocket::bind on the unspecified address UdpDialer::bind(AddressFamily::V4) or V6, with SocketOptions; the IPv6 socket is v6-only

Dialer also implements Connector<DialTarget>, which tries a list of TCP addresses in order or binds one socket per address family as a DualStackUdp. The dialers page covers both.

concepts/src/link.rs
pub enum NoStream {}
pub enum NoDatagram {
Never(Infallible),
}
impl AsyncRead for NoStream { /* match *self {} */ }
impl AsyncWrite for NoStream { /* match *self {} */ }
impl DatagramLink for NoDatagram {
type Addr = Destination;
// match *self { NoDatagram::Never(never) => match never {} }
}

These are the Stream and Datagram types of a connector that never produces that kind. Neither type can hold a value: NoStream has no variants, and NoDatagram’s one variant wraps an Infallible. Their trait methods are empty matches, which the compiler accepts because no value can reach them. An Outbound<TcpStream, NoDatagram> can only ever be Stream, so no check is needed at run time.

NoDatagram declares Addr = Destination deliberately. It then satisfies the runtime’s Conn::Datagram: DatagramLink<Addr = Destination> bound, so a stream-only connector can drive a ProxyServerRuntime.

Type Production user Test users
NoDatagram TransportConnector: a proxy’s UDP rides inside its stream, and a Udp destination is refused with a proxy transport carries no datagrams of its own the duplex dialers in concepts/tests/runtime.rs and concepts/tests/client.rs; protocols/tests/support/pipeline.rs → TcpConnector; katana’s tests/unit/serve.rs and tests/unit/e2e.rs
NoStream none the UDP connectors in concepts/tests/runtime.rs, typed Outbound<NoStream, UdpOutbound>

The codec-level counterpart is concepts/src/core.rs → NoCodec<Target, Error>, described on the client runtime page.

concepts/src/net.rs
pub enum UserAuthorization {
UsernamePassword {
username: CompactString,
password: CompactString,
},
Uuid(Uuid),
}

UserAuthorization is the identity a user was matched by. It is Clone, PartialEq and Eq, but not Hash. Protocol-specific key material, such as derived keys or cipher state, belongs in NetworkUser::user_data, not here. Every server core in etemenanki-protocols builds it once the user is authenticated, and none of them stores a secret in it:

Core Variant username password
SOCKS5 with accounts (socks/handshake.rs) UsernamePassword the matched account’s name empty
SOCKS5 without accounts, SOCKS4 UsernamePassword empty empty
HTTP (http/core.rs) UsernamePassword the matched account’s name, or empty when the inbound has no accounts empty
Trojan, Shadowsocks, Shadowsocks 2022 multi-user UsernamePassword the matched user’s email empty
Shadowsocks 2022 single key, TUN UsernamePassword empty empty
Hysteria 2 UsernamePassword the authenticated user’s label empty
VLESS, VMess Uuid

katana keys its per-user registry on this value. Because the type does not implement Hash, katana projects it onto its own AuthKey (Uuid or Name) in src/traffic.rs.

concepts/src/net.rs
pub struct NetworkUser<T> {
pub authorization: UserAuthorization,
pub user_data: std::sync::Arc<T>,
}
impl<T> Clone for NetworkUser<T> {
fn clone(&self) -> Self;
}

An authenticated user is its authorization plus an arbitrary per-user payload T behind an Arc. The payload is set by whoever built the inbound’s user table, and every flow the user opens shares the one instance:

Consumer T
etemenanki-app (app/src/flow.rs: Flow = etemenanki_protocols::flow::Flow<()>) ()
katana (src/connector.rs: Connector<Flow<UserTag>>) UserTag, the user’s AuthKey and panel uid; KatanaConnector::connect looks the user’s counter up from it per flow

Clone is written out rather than derived. #[derive(Clone)] would add a T: Clone bound, but the payload only ever lives behind an Arc, and Arc<T> clones for any T. A payload that must remain a single shared instance can therefore be a non-Clone type, and the user can still be cloned. The pipeline clones users often:

  • Flow::toward builds a sub-flow toward a new destination with the carrier’s user and source, and clones the user to do it. mux.cool sub-flows (protocols/src/mux/demux.rs), the flow a Trojan UDP association opens for its first packet and the app’s per-packet UDP routing (app/src/outbound/udp_fanout.rs) all go through it.
  • The SOCKS server clones the user into the one Flow it dials for a UDP association, when the first datagram from its client is forwarded. That datagram’s target becomes the flow’s destination. The Shadowsocks 2022 and VMess cores clone it into each request’s Flow.

Flow<T> writes its own Clone for the same reason, and so does UserEntry<T> in protocols/src/hysteria/server/authenticator.rs.

concepts/src/net.rs
#[repr(u8)]
pub enum DialNetwork {
Unknown = 0,
Tcp = 1,
Udp = 2,
Unix = 3,
}

DialNetwork is the transport network of a Destination. It is Copy, Eq and Hash. The protocol cores build only Tcp and Udp, and every consumer branches on those two:

  • AppConnector, KatanaConnector and FreedomConnector take the datagram path for Udp.
  • TransportConnector refuses Udp with ErrorKind::Unsupported.
  • route_target, in both app/src/router.rs and katana’s src/router.rs, maps Tcp and Udp to the route model’s TargetNetwork and leaves the network unset for any other value.
concepts/src/net.rs
pub enum Remote {
IpAddr(std::net::IpAddr),
Domain(CompactString),
}

Remote is a host that has not necessarily been resolved. A domain read from the wire stays a domain until it reaches the component that must contact it. Routing can then match it by name, and the chosen outbound’s resolver and address-family strategy decide how to reach it. No server core resolves a name.

The route model has its own borrowed Remote<'a> in environment/src/routing.rs. route_target converts one into the other, borrowing the domain string rather than copying it.

concepts/src/net.rs
pub struct Destination {
pub network: DialNetwork,
pub remote: Remote,
pub port: u16,
}
impl Destination {
pub fn udp(addr: std::net::SocketAddr) -> Self;
pub fn socket_addr(&self) -> Option<std::net::SocketAddr>;
}

Destination is what a request names. It is also the per-packet address of a datagram through an outbound: Effect::SendTo, Effect::SendToHeld, Event::Datagram and Event::SendFailed all carry one. Every UDP-carrying protocol puts a host that may be a domain on the wire for each packet, so the core passes that host on unchanged and the link decides how to reach it. Destination is Clone, Eq and Hash.

Helper Behaviour
Destination::udp(addr) network: DialNetwork::Udp, remote: Remote::IpAddr(addr.ip()), port: addr.port(). The plain socket links use it to report the source of a received packet.
destination.socket_addr() Some(SocketAddr) for Remote::IpAddr, None for Remote::Domain. It ignores network. The plain UDP outbounds use it to refuse domains.
concepts/src/sniff.rs
pub enum SniffedProtocol {
Tls,
Http,
}
pub struct SniffedBehavior {
pub protocol: SniffedProtocol,
pub domain: CompactString,
}
pub trait Sniffer {
fn sniff(&self, data: &[u8]) -> Option<SniffedBehavior>;
}

A Sniffer inspects the leading payload bytes of a flow and may recover a routing hint the proxy header did not carry: a TLS SNI or an HTTP Host. sniff takes &self and a plain slice and returns an Option, so a sniffer has no way to fail a connection. Malformed, truncated or unrecognised input yields None, and the flow routes on its destination as before.

The concepts crate defines only these shapes. protocols/src/sniff/ provides the rest:

  • TlsSniffer and HttpSniffer, the two Sniffer implementations;
  • the function sniff(data), which tries TLS first because a TLS record header is the tighter discriminator;
  • plausible_domain, which rejects IP literals, names longer than 253 bytes and characters outside [A-Za-z0-9._-];
  • worth_sniffing(destination), which is true only for a Remote::IpAddr destination;
  • the budgets SNIFF_TIMEOUT (300 ms) and SNIFF_LIMIT (4 KiB, counted across frames).

The result travels in Flow::sniffed. An inbound sets it before it dials, and only when the inbound has sniffing enabled and worth_sniffing allowed it. Routing reads only domain: both route_target functions pass it to RouteTarget::with_sniffed_domain, and every domain matcher then tries it as well as the request’s own host. protocol is carried but no consumer reads it. The sniffing page covers the collector and each protocol’s hold-and-open logic.

relay.rs holds two hand-written futures for paths that have no protocol between the two sides. Each keeps one fixed buffer per direction, heap-allocated by boxed_array at construction. There is no task per direction, no channel and no allocation after new. Neither future has a timeout of its own; a caller that needs one wraps the future.

concepts/src/relay.rs
pub struct UnidirectionalConnection<I, O, const BUF_SIZE: usize = 8192> { /* … */ }
impl<I, O, const BUF_SIZE: usize> UnidirectionalConnection<I, O, BUF_SIZE> {
pub fn new(input: I, output: O) -> Self;
pub fn transferred(&self) -> u64;
}
impl<I: AsyncRead, O: AsyncWrite, const BUF_SIZE: usize> Future
for UnidirectionalConnection<I, O, BUF_SIZE>
{
type Output = io::Result<u64>;
}

UnidirectionalConnection copies input to output until input reaches end of stream, flushes output, and resolves to the byte count. input and output are structurally pinned with #[pin_project], so neither has to be Unpin. At end of stream it flushes but does not shut output down. No code in the workspace or in katana uses it at the pinned revisions.

concepts/src/relay.rs
pub struct Relayed {
pub a_to_b: u64,
pub b_to_a: u64,
}
pub struct BidirectionalConnection<A, B, const BUF_SIZE: usize = 8192> { /* … */ }
impl<A, B, const BUF_SIZE: usize> BidirectionalConnection<A, B, BUF_SIZE> {
pub fn new(a: A, b: B) -> Self;
pub fn relayed(&self) -> Relayed;
}
impl<A, B, const BUF_SIZE: usize> Future for BidirectionalConnection<A, B, BUF_SIZE>
where
A: AsyncRead + AsyncWrite + Unpin,
B: AsyncRead + AsyncWrite + Unpin,
{
type Output = io::Result<Relayed>;
}

BidirectionalConnection copies a → b and b → a in one future. Each direction is a private Half<BUF_SIZE> with its own buffer, cursors (pos, cap), flags (read_done, shutdown_done, need_flush) and byte count. Every poll steps both halves with the same Context, lending the same endpoint to one half as its reader and to the other as its writer. A socket therefore needs no split. The endpoints are not pinned, which is why the Future impl requires Unpin on both. The future resolves once both halves are done. Relayed is Copy and PartialEq, and relayed() can be read at any time, which is how a caller detects an idle relay.

One half moves through these states:

stateDiagram-v2
  [*] --> Fill
  Fill --> Drain: read returned bytes
  Fill --> Parked: read Pending (flush first if need_flush)
  Parked --> Fill: woken
  Fill --> ShutDown: read returned 0 bytes
  Drain --> Drain: partial write
  Drain --> Fill: buffer drained
  ShutDown --> Done: poll_shutdown ready
  Fill --> Failed: read error
  Drain --> Failed: write error or zero-length write
  ShutDown --> Failed: shutdown error
  Done --> [*]
  Failed --> [*]
  • Parked. Before a half returns Pending on an idle reader, it flushes the writer if it wrote anything since the last flush (need_flush). Bytes buffered in a TLS or WebSocket writer therefore reach the peer even when no more input arrives.
  • ShutDown. At end of stream the half calls poll_shutdown on its writer: a half-close. The other direction keeps flowing. UnidirectionalConnection flushes at this point instead.
  • Failed. A read error, a write error, or a write that returns 0 (ErrorKind::WriteZero) ends the whole future with that error. The ? after the a → b half returns before the b → a half is polled in that call.

The protocol cores run inside the server runtime and do not use relay.rs. At 596916d its only caller is the SOCKS CONNECT path in protocols/src/socks/server.rs. SOCKS is the one inbound whose server does not fit a ProxyCoreDecode, because its UDP side lives on a second socket. After the handshake and the connector’s dial, the CONNECT path relays the client stream and the connector’s stream through:

protocols/src/socks/server.rs
async fn relay_with_idle_guard<A, B>(a: A, b: B) -> io::Result<Relayed>
where
A: AsyncRead + AsyncWrite + Unpin,
B: AsyncRead + AsyncWrite + Unpin,

It builds a BidirectionalConnection::<A, B, RELAY_BUF> (16 KiB per direction) and, in a loop, races it against a fresh RELAY_IDLE_TIMEOUT (300 s) sleep. When the sleep fires, it compares relayed() with the previous snapshot. If nothing moved, it returns ErrorKind::TimedOut with the text socks: relay idle. Otherwise it stores the new snapshot and sleeps again. A relay that falls silent is therefore closed between one and two RELAY_IDLE_TIMEOUT periods later. Returning drops the future, which drops both endpoints and closes both connections.

The UDP ASSOCIATE path does not use relay.rs either. It drives the connector’s DatagramLink from its own select! loop over the hub socket, outside any runtime, and ExpectedSender in protocols/src/socks/server.rs decides which hub datagrams reach that link:

  • admits(from) requires the control connection’s IP, or over a Unix socket the IP the request named. It compares with endpoint from protocols/src/socks/protocol.rs, which puts the IP in canonical form, so an IPv4-mapped IPv6 sender is the IPv4 one. Once a port is fixed, by the request or by pin, it requires that port too. A datagram that fails the check is dropped before it is parsed.
  • pin(from) runs for every admitted datagram that parses and carries a non-empty payload, before the link is dialed and the payload is sent. The first pin holds, so the first such datagram fixes the client, and client() then names the one address that replies from the link go to.
  • The loop dials the association’s single link for that same datagram. A datagram from any other sender therefore never opens an outbound and never receives a reply.

The SOCKS page covers the port a request may pin up front, the Unix-socket case and the 0x02 refusal of an association the relay could not hold to its client.

The client-side link uses the same comparison. SocksUdpLink::poll_recv_from in protocols/src/socks/udp_link.rs skips every datagram whose endpoint differs from the relay’s, as well as one that does not parse. A link whose socket is bound dual-stack on [::] therefore still hears an IPv4 relay, which that socket reports in IPv4-mapped form.

katana has no SOCKS inbound and does not reference relay.rs.

Section titled “Data flow: one UDP packet through an outbound link”

The datagram path shows the link contract most clearly. The core, the runtime, the connector and the link all run in one task, and every arrow below is a function call or a poll, never a channel.

sequenceDiagram
  participant C as ProxyCoreDecode
  participant R as ProxyServerRuntime
  participant K as Connector
  participant L as DatagramLink (Addr = Destination)
  C->>R: Effect::Open (key, target: Flow)
  R->>K: connect(target), future boxed in the slot
  K-->>R: Ok(Outbound::Datagram(link))
  R->>C: Event::Connected (key)
  C->>R: Effect::SendTo (key, to: Destination, range)
  R->>L: poll_send_to(bytes, to)
  alt Ready(Ok(n))
    R->>R: outbound_tx += n
  else Ready(Err(e)), for example a domain on a plain link
    R->>C: Event::SendFailed (key, to, error), key stays live
  end
  L-->>R: poll_recv_from returns Ok(from)
  R->>C: Event::Datagram (key, from, data)
  L-->>R: poll_recv_from returns Err(e)
  R->>C: Event::OutboundError (key, error), key removed

A stream outbound follows the same shape with Effect::Forward and Effect::Shutdown, the AsyncRead and AsyncWrite poll methods, and Event::Outbound and Event::OutboundEof. On a stream outbound, any read or write error removes the key. The connection lifecycle page follows a whole connection end to end.

Invariant Mechanism Pinned by
A live outbound can be moved. DatagramLink: Unpin and Connector::Stream: Unpin. The runtime keeps slots by value in a BTreeMap. The compiler
The server runtime only sends datagrams addressed by Destination. Conn::Datagram: DatagramLink<Addr = Destination> on the runtime’s impls. NoDatagram declares that Addr so stream-only connectors qualify. The compiler
A plain UDP outbound never resolves a name. UdpOutbound and DualStackUdp call socket_addr() and return ErrorKind::Unsupported on None. No test sends a domain through either link. The refusal path they share is pinned by a_refused_datagram_send_keeps_the_key_alive in concepts/tests/runtime.rs, which sends to an IPv6 peer through an IPv4 socket.
A refused datagram send keeps the key; a failed receive removes it. report_send_failure in concepts/src/runtime.rs leaves the slot alone; fail_outbound calls forget_key. a_refused_datagram_send_keeps_the_key_alive in concepts/tests/runtime.rs
A refused packet to a datagram transport peer is not fatal. Send errors on a datagram transport become Event::TransportSendFailed. a_refused_transport_packet_is_reported_not_fatal and single_peer_datagram_transport_demuxes_sessions_and_reports_refusals in concepts/tests/runtime.rs
A failed dial reaches the core as an event, not as a runtime error. Poll::Ready(Err(error)) from the boxed dial future → forget_key → Event::ConnectFailed. connect_failure_reaches_the_core_as_an_event in concepts/tests/runtime.rs
A stream effect never lands on a datagram outbound, or the reverse. The LinkState match arms in the effect loop return RuntimeError::WrongLinkKind. No dedicated test
A client runtime’s wire is always a stream. ProxyClientRuntime::wire maps Outbound::Datagram to ErrorKind::Unsupported and sets Wire::Down. No dedicated test
Cloning a user never requires T: Clone. The hand-written impl<T> Clone for NetworkUser<T> and impl<T> Clone for Flow<T> clone the Arc. The compiler, for every payload type in use; no test uses a non-Clone payload
A domain from the wire reaches routing unresolved. Remote::Domain is carried in Destination, and route_target maps it to routing::Remote::Domain. On the core side, connect_to_a_domain_answers_200_once_connected in protocols/tests/unit/http/core.rs checks that the opened Flow names the domain. No test pins the route_target conversion itself.
A sniffed domain lets a domain rule match an IP-addressed flow. Flow::sniffed, read by route_target through RouteTarget::with_sniffed_domain an_http_host_routes_an_ip_addressed_flow, a_tls_sni_routes_an_ip_addressed_flow and turning_sniffing_off_stops_the_domain_rule_matching in app/tests/integration/e2e_sniff.rs; a_sniffed_domain_makes_an_ip_target_match_domain_rules in environment/tests/unit/routing.rs
A relay half-closes each direction on its own and reports exact counts. Half::poll calls poll_shutdown at end of stream, and the future resolves only when both halves are Ready. bidirectional_relays_both_ways_and_half_closes in concepts/src/relay.rs, with a 16-byte buffer and a 41-byte message so the copy wraps the buffer
Buffered relay output is not stranded when input goes quiet. need_flush is flushed before a half returns Pending on its reader. No dedicated test
A SOCKS UDP association feeds its link only from its own client, and replies only to it. ExpectedSender::admits runs before the packet is parsed; pin runs before the link is dialed; replies go to client(). only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port and an_ipv4_mapped_address_is_the_ipv4_one in protocols/tests/unit/socks/server.rs; udp_association_ignores_another_ip and udp_association_ignores_another_port_once_pinned in protocols/tests/pipeline/socks.rs
A SocksUdpLink returns only the relay’s datagrams, whatever form its socket reports the relay’s address in. poll_recv_from compares endpoint(from) with endpoint(self.relay) and skips a mismatch. udp_link_ignores_datagrams_not_from_the_relay and udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay in protocols/tests/pipeline/socks.rs; endpoint_sees_through_ipv4_mapping_and_ignores_flow_info in protocols/tests/unit/socks/protocol.rs
Where Failure Outcome
Connector::connect future Err(e) Server runtime: the key is forgotten and the core gets Event::ConnectFailed. Client runtime: Wire::Down, the first plaintext call fails with e, and every later call fails with NotConnected. SOCKS CONNECT: a refusal reply (SOCKS5 code 0x05 for ConnectionRefused, 0x04 for anything else). When sniffing already granted the request and collected a prefix of client bytes, no refusal is written and the client sees the connection close.
DatagramLink::poll_send_to, outbound Err(e) The packet is dropped, the core gets Event::SendFailed, and the key stays.
DatagramLink::poll_recv_from, outbound Err(e) The key is removed and the core gets Event::OutboundError. A ProxyClientRuntime datagram link returns UnexpectedEof here once the upstream has ended the flow.
DatagramLink, transport send Err / receive Err Event::TransportSendFailed / RuntimeError::Transport, which ends the runtime.
Wrong Outbound variant see the table under Outbound<S, D> RuntimeError::WrongLinkKind, or ErrorKind::Unsupported from the consumer.
BidirectionalConnection read, write or shutdown error, or WriteZero, in either half The future resolves to Err. The caller drops it, which drops both endpoints.
SOCKS relay no bytes in either direction across a full RELAY_IDLE_TIMEOUT check interval ErrorKind::TimedOut, socks: relay idle.

Cancellation is drop. Every type on this page is owned by the future that polls it. Dropping a ProxyServerRuntime drops the dial futures still in flight and every link in its map. Dropping a relay future drops the endpoints it owns. None of these types spawns a task, so nothing they own outlives the connection.

Constant Value Defined in Applies to
BUF_SIZE default of both relay futures 8192 bytes per direction concepts/src/relay.rs UnidirectionalConnection, BidirectionalConnection
RELAY_BUF 16 * 1024 bytes per direction protocols/src/socks/server.rs The SOCKS CONNECT relay
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs The SOCKS relay’s idle guard, and the cores’ idle deadline
ProxyCoreDecode::MAX_DATAGRAM 4096 bytes by default, overridable per core concepts/src/core.rs The largest packet the server runtime reads whole from a DatagramLink; capped at BUF_SIZE - STAGING_RESERVE for an outbound and at BUF_SIZE for a transport
DEFAULT_CONNECT_TIMEOUT 10 s environment/src/dial/tcp.rs Dialer’s TCP connects, unless TcpDialer::with_connect_timeout overrides it; SocketConnector has none
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs How long an inbound waits for the first payload before routing without a sniffed domain
SNIFF_LIMIT 4 * 1024 bytes protocols/src/sniff/mod.rs Bytes a sniffer may inspect, across frames
Test File What it pins
bidirectional_relays_both_ways_and_half_closes concepts/src/relay.rs Both directions, half-close, and Relayed { a_to_b: 41, b_to_a: 5 } with a buffer smaller than the message
datagrams_round_trip_through_a_real_udp_socket concepts/tests/runtime.rs A closure connector returning Outbound::<NoStream, UdpOutbound>::Datagram under a stream transport, with byte counts
a_refused_datagram_send_keeps_the_key_alive concepts/tests/runtime.rs A poll_send_to error becomes SendFailed, and the same key relays afterwards
datagram_outbound_is_never_truncated_by_staging_backpressure concepts/tests/runtime.rs A datagram outbound is read only with room for a whole packet
connect_failure_reaches_the_core_as_an_event concepts/tests/runtime.rs A closure connector’s Err becomes ConnectFailed
datagram_transport_demultiplexes_peers_and_frames_replies concepts/tests/runtime.rs A bare UdpSocket (Addr = SocketAddr) as a datagram transport, with the function udp_connector as the connector
a_refused_transport_packet_is_reported_not_fatal concepts/tests/runtime.rs A transport-side send error becomes TransportSendFailed
single_peer_datagram_transport_demuxes_sessions_and_reports_refusals concepts/tests/runtime.rs A link with Addr = () (TinyQuic) as a transport
dial_failure_surfaces_on_first_use concepts/tests/client.rs A closure connector’s dial error reaches the client runtime’s first call, then NotConnected on later calls
socket_target_connector_binds_one_family environment/tests/integration/udp.rs Dialer as Connector<SocketTarget> binds one IPv4 socket and yields a UdpOutbound
dual_stack_link_serves_a_proxy_runtime environment/tests/integration/udp.rs Dialer passed directly as a server runtime’s connector, yielding a DualStackUdp that relays a packet both ways
transport_connector_refuses_udp protocols/tests/pipeline/transports.rs TransportConnector::dial, which connect wraps, refuses a Udp destination with ErrorKind::Unsupported
new_server_vs_new_client_tcp, new_server_vs_new_client_udp protocols/tests/pipeline/socks.rs The SOCKS CONNECT relay end to end, and SocksUdpLink as a DatagramLink<Addr = Destination>
udp_association_ignores_another_ip (Linux only), udp_association_ignores_another_port_once_pinned protocols/tests/pipeline/socks.rs A datagram from another IP, or from another port once the client is pinned, is never forwarded, and the sender on another IP gets no reply
only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port, an_ipv4_mapped_address_is_the_ipv4_one protocols/tests/unit/socks/server.rs ExpectedSender: the control peer’s IP only, the first pin holds, and an IPv4-mapped sender matches its IPv4 peer
udp_link_ignores_datagrams_not_from_the_relay protocols/tests/pipeline/socks.rs SocksUdpLink::poll_recv_from skips an intruder’s datagram and returns the relay’s
udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay (Linux only) protocols/tests/pipeline/socks.rs A SocksUdpLink on a [::] socket relays through an IPv4 relay
endpoint_sees_through_ipv4_mapping_and_ignores_flow_info protocols/tests/unit/socks/protocol.rs endpoint equates IPv4-mapped and IPv4 addresses and ignores IPv6 flow info
extracts_sni_from_a_real_client_hello, a_truncated_client_hello_yields_nothing_rather_than_garbage, an_ip_literal_sni_is_rejected, non_tls_and_malformed_input_yield_nothing protocols/tests/unit/sniff/tls.rs TlsSniffer as a Sniffer
extracts_the_host_header, an_ip_host_is_rejected, bytes_that_merely_contain_a_host_header_are_not_http protocols/tests/unit/sniff/http.rs HttpSniffer as a Sniffer
domain_plausibility protocols/tests/unit/sniff/mod.rs plausible_domain, the filter every sniffed name passes through

UnidirectionalConnection and SocketConnector have no test and no caller.