Skip to content

Outbounds, fan-out and balancing

Source files: 28 · checked against Etemenanki 596916d
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/outbound/proxy.rs
  • Etemenanki/app/src/outbound/freedom.rs
  • Etemenanki/app/src/outbound/udp_fanout.rs
  • Etemenanki/app/src/connector.rs
  • Etemenanki/app/src/balancer.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/router.rs
  • Etemenanki/app/src/flow.rs
  • Etemenanki/app/src/config.rs
  • Etemenanki/concepts/src/link.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/socks/udp_link.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/src/transports/connect.rs
  • Etemenanki/protocols/src/vless/codec.rs
  • Etemenanki/protocols/src/vmess/codec.rs
  • Etemenanki/environment/src/dial/tcp.rs
  • Etemenanki/environment/src/dial/udp.rs
  • Etemenanki/app/tests/unit/outbound.rs
  • Etemenanki/app/tests/unit/balancer.rs
  • Etemenanki/app/tests/integration/e2e_udp_route.rs
  • Etemenanki/app/tests/integration/e2e_balancer.rs
  • Etemenanki/app/tests/integration/e2e_route_context.rs
  • Etemenanki/protocols/tests/pipeline/socks.rs
  • Etemenanki/app/tests/integration/e2e_xray.rs

Every inbound in etemenanki-app ends in the same place: a server core asks its runtime to open a flow, and the runtime calls a Connector. In the app that connector is AppConnector. It routes the flow, picks an Outbound, and hands back either a byte stream or a datagram link. This page follows that path from the connector down to the sockets: the Outbound enum and its two dispatch methods, the per-protocol clients, the direct (freedom) connector, the per-packet UDP fan-out, and the balancer that picks among several upstreams by health.

It is written for contributors who add an outbound protocol, change how UDP associations are routed, or touch balancing. It assumes you know the server runtime from The server runtime and the client side from Client codecs and the client runtime. The operator’s view of the same features is in Outbounds and Balancers.

Component Owns Leaves to others
AppConnector (app/src/connector.rs) Routing a TCP flow once and dialing it; turning a UDP flow into a FanOutLink instead of dialing it. The route table (Router), the dial itself (Outbound).
Outbound (app/src/outbound/mod.rs) One handler per configured outbound: which client opens a flow, and whether it can carry TCP, UDP or both. Protocol framing (the codecs in etemenanki-protocols), transports (TransportConnector).
ProxyClient (app/src/outbound/proxy.rs) Sharing one ProxyClientConnector between all flows of an outbound; choosing the stream or datagram codec per flow. Driving the codec over the wire (ProxyClientRuntime).
SocksOutbound (app/src/outbound/mod.rs) SOCKS5 CONNECT through ProxyClient, and UDP ASSOCIATE as its own link. The association’s wire format (SocksUdpLink in protocols/src/socks/udp_link.rs).
FreedomConnector and ResolvingUdp (app/src/outbound/freedom.rs) Dialing the real destination: TCP resolve and connect, UDP bind per family and per-name resolution. Socket policy (Dialer), DNS (Resolver).
FanOutLink (app/src/outbound/udp_fanout.rs) Routing every packet of a UDP association on its own, keeping one sub-link per routed outbound, merging replies. Queueing and backpressure (the server runtime keeps a waiting packet queued).
Balancer (app/src/balancer.rs) Health state per member, one probe task per member, selection per flow. Building members and deciding what may be balanced (build in app/src/instance.rs).

Nothing in this layer keeps connections alive across flows except what the protocol clients themselves share (the WireGuard device and the Hysteria 2 connection). Everything else is created per flow or per association and dropped with it.

build in app/src/instance.rs turns the config into one generation’s worth of objects:

  1. build_outbound runs once per [[outbound]]. Each result goes into a HashMap<CompactString, Arc<Outbound>> keyed by tag. A duplicate tag fails with duplicate outbound tag: <tag>.
  2. Each [[balancer]] becomes an Outbound::Balanced(Arc<Balancer>) in the same map, so a route rule can name a balancer wherever it can name an outbound.
  3. build_router resolves every rule’s tag against the map. The Router (routing::Router<Outbound>) holds the Arc<Outbound>s for as long as the generation lives.
  4. spawn_generation starts the balancer probes under the generation’s CancellationToken, then binds the inbounds. Each inbound task builds an AppConnector from the shared Arc<Router> and a FlowContext.

--test runs build and stops, so it validates every outbound and balancer without starting a probe or binding a listener.

app/src/connector.rs
#[derive(Clone)]
pub struct AppConnector {
pub router: Arc<Router>,
pub ctx: FlowContext,
}
type ConnectFuture =
Pin<Box<dyn Future<Output = io::Result<Outbound<OutboundStream, FanOutLink>>> + Send>>;
impl Connector<Flow> for AppConnector {
type Stream = OutboundStream;
type Datagram = FanOutLink;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow) -> ConnectFuture;
}

Outbound in the ConnectFuture is etemenanki_concepts::link::Outbound, the two-armed “stream or datagram” result, not the app’s Outbound enum. Flow is etemenanki_protocols::flow::Flow<()>: a destination, the NetworkUser (with a unit payload in the app), what sniffing recovered (sniffed, optional) and an optional source address. FlowContext carries the inbound tag and the listener-side peer address, which the route matchers inbound_tag and source_cidr need.

connect branches on flow.destination.network:

  • UDP. No route lookup, no dial. It builds a FanOutLink over the router, the context and the flow, and returns a future that is ready at once. The server core therefore sees Connected for a UDP flow immediately, and every later packet is routed on its own.
  • Anything else. It calls router.route(route_target(&flow, &self.ctx)) synchronously, calls connect_stream on the chosen outbound, and maps the result into Outbound::Stream. Routing and balancer selection both happen when connect is called, before the future is first polled.

route_target (app/src/router.rs) maps the flow to a routing::RouteTarget: the destination, the sniffed domain if there is one, the inbound tag, flow.source.or(ctx.source) as the source, and the network (Tcp or Udp) from DialNetwork.

app/src/outbound/mod.rs
pub enum Outbound {
Freedom(FreedomConnector),
Blackhole,
Socks(Box<SocksOutbound>),
Http(ProxyClient<HTTP_BUF, HttpConnect, NoUdp>),
Trojan(ProxyClient<TROJAN_BUF, TrojanStream, TrojanDatagram>),
Vless(ProxyClient<VLESS_BUF, VlessStream, VlessDatagram>),
Vmess(ProxyClient<VMESS_BUF, VMessStream, VMessDatagram>),
Shadowsocks(ProxyClient<SS_BUF, SsStream, NoUdp>),
Ss2022(ProxyClient<SS2022_BUF, Ss2022Stream, NoUdp>),
Wireguard(WgConnector),
Hysteria2(Hy2Connector),
Balanced(Arc<crate::balancer::Balancer>),
}
pub type StreamFuture = Pin<Box<dyn Future<Output = io::Result<OutboundStream>> + Send>>;
pub type DatagramFuture = Pin<Box<dyn Future<Output = io::Result<OutboundDatagram>> + Send>>;
impl Outbound {
pub fn connect_stream(&self, flow: Flow) -> StreamFuture;
pub fn connect_datagram(&self, flow: Flow) -> DatagramFuture;
}
pub fn build_outbound(cfg: &OutboundConfig, resolver: &Resolver) -> io::Result<Outbound>;

Both dispatch methods take &self, because the router hands out shared Arc<Outbound>s and many flows dial the same outbound at once. Each arm builds its dial future synchronously and boxes it. The table shows what each arm does:

Variant connect_stream connect_datagram
Freedom Clones the connector, dials, expects link::Outbound::Stream → OutboundStream::Tcp Clones the connector, dials, expects link::Outbound::Datagram → OutboundDatagram over ResolvingUdp
Blackhole Ready at once: OutboundStream::Blackhole Ready at once: OutboundDatagram over BlackholeLink
Socks SocksOutbound.tcp.connect(flow) → proxy_stream SocksOutbound::connect_datagram(), a UDP ASSOCIATE (the flow is not used)
Http ProxyClient::connect → proxy_stream Fails: http carries no datagrams
Trojan, Vless, Vmess ProxyClient::connect → proxy_stream ProxyClient::connect → proxy_datagram
Shadowsocks, Ss2022 ProxyClient::connect → proxy_stream Fails: shadowsocks carries no datagrams, shadowsocks-2022 carries no datagrams
Wireguard Clones the connector, dials, expects a stream → OutboundStream::Wg Clones the connector, dials, expects a datagram link → OutboundDatagram
Hysteria2 Clones the connector, dials, expects a stream → OutboundStream::Hy2 Clones the connector, dials, expects a datagram link → OutboundDatagram
Balanced b.select().connect_stream(flow) b.select().connect_datagram(flow)

The “no datagrams” errors have kind Unsupported. WgConnector and Hy2Connector are cheap to clone: their clones share the tunnel device or QUIC connection through an Arc<Mutex<…>> slot, so cloning per dial does not open a new tunnel.

Seven protocols are clients over a transport: HTTP, SOCKS CONNECT, Trojan, VLESS, VMess, Shadowsocks and Shadowsocks 2022. All of them go through the same wrapper:

app/src/outbound/proxy.rs
pub type NoUdp = NoCodec<Destination, io::Error>;
pub type Make<S, D> = Box<dyn FnMut(Flow) -> link::Outbound<S, D> + Send>;
pub struct ProxyClient<const BUF: usize, S, D> {
inner: Mutex<ProxyClientConnector<BUF, Make<S, D>, TransportConnector, Destination>>,
}
impl<const BUF: usize, S, D> ProxyClient<BUF, S, D>
where
S: ProxyCoreEncode<Target = Destination, Error = io::Error>,
D: ProxyCoreEncodeDatagram<Target = Destination, Error = io::Error>,
{
pub fn new(make: Make<S, D>, transport: TransportConnector, server: Destination) -> Self;
pub fn connect(&self, flow: Flow) -> ProxyClientConnecting<BUF, S, D, TransportConnector>;
}

Why a mutex. ProxyClientConnector::connect takes &mut self: Make is an FnMut, and the inner TransportConnector::connect takes &mut self as well. The outbound is shared behind an Arc, so ProxyClient puts the connector behind a parking_lot::Mutex. connect locks, calls the inner connect, and returns. Under the lock the connector only runs make, builds the ProxyClientRuntime (codec plus buffers) and asks TransportConnector for its dial future, which clones the transport into a boxed async block. No I/O and no .await happen while the lock is held, so flows through one outbound serialise only on building a future, never on the network.

The codec is picked per flow. build_outbound gives each protocol a Make closure that captures the credentials (a password hash, a UUID, a derived key) and builds a codec for the flow’s destination:

Protocol Stream codec Datagram codec How make chooses
http HttpConnect::new(&flow.destination, auth) NoUdp Always a stream
socks SocksConnect::new(&flow.destination, auth) NoUdp Always a stream; UDP bypasses ProxyClient
trojan TrojanStream::new(&hash, &flow.destination) TrojanDatagram::new(&hash) flow.destination.network == DialNetwork::Udp
vless VlessStream::new(&uuid, &flow.destination) VlessDatagram::new(&uuid, &flow.destination) same
vmess VMessStream::new(uuid, security, true, &flow.destination) VMessDatagram::new(uuid, security, true, &flow.destination) same; true is global_padding
shadowsocks (legacy methods) SsStream::new(method, key.clone(), &flow.destination) NoUdp Always a stream
shadowsocks (2022- methods) Ss2022Stream::new(method, psk.clone(), keys.clone(), &flow.destination) NoUdp Always a stream

A NoUdp protocol never reaches its make with a UDP flow: connect_datagram refuses it first.

The dial future is a ProxyClientConnecting. It resolves only once the upstream is dialed, the codec’s handshake is done and the handshake bytes are flushed, so a refused upstream is an Err from the future rather than a stream that fails later. Two adapters then fold the result into the app’s closed types:

app/src/outbound/proxy.rs
pub fn proxy_stream<const BUF: usize, S, D>(
opened: link::Outbound<
ProxyClientRuntime<BUF, S, TransportConnector>,
ProxyClientRuntime<BUF, D, TransportConnector>,
>,
) -> io::Result<OutboundStream>
where
S: ProxyCoreEncode<Target = Destination, Error = io::Error> + Send + 'static,
D: ProxyCoreEncodeDatagram<Target = Destination, Error = io::Error> + Send + 'static;
pub fn proxy_datagram<const BUF: usize, S, D>(
opened: link::Outbound<
ProxyClientRuntime<BUF, S, TransportConnector>,
ProxyClientRuntime<BUF, D, TransportConnector>,
>,
) -> io::Result<OutboundDatagram>
where
S: ProxyCoreEncode<Target = Destination, Error = io::Error> + Send + 'static,
D: ProxyCoreEncodeDatagram<Target = Destination, Error = io::Error> + Send + 'static;

A runtime of the wrong kind is an error: a TCP flow was dialed as datagrams or a UDP flow was dialed as a stream. Neither happens in practice, because make decides the kind from the same network field that chose between connect_stream and connect_datagram.

Section titled “OutboundStream, OutboundDatagram and BlackholeLink”

The server runtime stores outbounds by value and needs one concrete type for each kind, so the app closes them:

app/src/outbound/proxy.rs
pub trait Stream: AsyncRead + AsyncWrite + Send {}
pub type ProxyStream = Pin<Box<dyn Stream>>;
pub enum OutboundStream {
Tcp(TcpStream),
Proxy(ProxyStream),
Wg(WgStream),
Hy2(Hy2Stream),
Blackhole,
}
pub struct OutboundDatagram(Box<dyn DatagramLink<Addr = Destination> + Send>);
impl OutboundDatagram {
pub fn new<D: DatagramLink<Addr = Destination> + Send + 'static>(link: D) -> Self;
}
pub struct BlackholeLink;
  • OutboundStream implements AsyncRead and AsyncWrite by delegating to the variant. The proxy clients differ in type per protocol and per buffer size, so they are boxed into Proxy. Plain TCP, WireGuard and Hysteria 2 streams keep their own variants and avoid the box.
  • OutboundStream::Blackhole reads EOF at once and accepts every write whole. flush and shutdown succeed immediately.
  • OutboundDatagram boxes every datagram link, because the fan-out holds links of different outbounds side by side.
  • BlackholeLink reports every poll_send_to as sent in full and returns Pending from every poll_recv_from, forever.
app/src/outbound/mod.rs
pub struct SocksOutbound {
transport: TransportConnector,
server: Destination,
auth: Option<(CompactString, CompactString)>,
tcp: ProxyClient<SOCKS_BUF, SocksConnect, NoUdp>,
}

TCP flows go through tcp, an ordinary ProxyClient. UDP does not fit the client runtime, because a SOCKS5 association is two connections: a control stream that must stay open, and a separate UDP socket toward the server’s relay. connect_datagram therefore builds its own link:

  1. Dials a control stream with transport.dial(&server), so the [outbound.stream] transport (TLS, WebSocket, gRPC) applies to the control stream.
  2. Runs SocksUdpLink::associate(control, auth, bind): the method negotiation, the username and password round if auth is set, then UDP ASSOCIATE.
  3. Binds a plain UDP socket through the bind closure, on 0.0.0.0:0 or [::]:0 to match the relay address’s family.

The UDP ASSOCIATE request names no source: encode_request(CMD_UDP_ASSOCIATE, None) writes 0.0.0.0:0, as RFC 1928 has a client do when it does not know its source. The socket is bound only after the reply names the relay, and a local address would be the wrong one behind NAT. An upstream that checks sources, as the Etemenanki SOCKS inbound does, therefore holds the association to the IP the control connection came from and the port of the first datagram it relays. The datagrams must leave from the IP the upstream sees on the control connection. When the control stream reaches the upstream through a reverse proxy or a CDN (for example a WebSocket or gRPC transport behind one), the upstream sees that front end’s address instead, and a source-checking upstream drops the association’s datagrams.

protocols/src/socks/udp_link.rs
impl<S> SocksUdpLink<S>
where
S: AsyncRead + AsyncWrite + Unpin,
{
pub async fn associate(
mut control: S,
auth: Option<(&str, &str)>,
bind: impl FnOnce(&SocketAddr) -> io::Result<UdpSocket>,
) -> io::Result<Self>;
}
sequenceDiagram
  participant F as FanOutLink
  participant O as SocksOutbound
  participant T as TransportConnector
  participant S as SOCKS5 server
  F->>O: connect_datagram
  O->>T: dial server
  T->>S: TCP connect, then transport handshake
  O->>S: method request
  S-->>O: chosen method
  opt auth is set
    O->>S: username and password
    S-->>O: status
  end
  O->>S: UDP ASSOCIATE naming source 0.0.0.0 port 0
  S-->>O: reply with relay address
  O->>O: bind UDP socket in the relay's family
  O-->>F: OutboundDatagram over SocksUdpLink
  F->>O: poll_send_to
  O->>S: relay header plus payload, to the relay address
  S-->>O: datagram from the relay address
  O-->>F: payload and the peer named in its header

Once associated, the link behaves as follows:

  • Every poll_send_to wraps the payload in the SOCKS5 UDP header for that packet’s own destination, so one association reaches many peers.
  • poll_recv_from drops any datagram whose source is not the relay, and any datagram whose header does not parse. It then returns the peer the header names. The source check is endpoint(from) != endpoint(self.relay). endpoint (protocols/src/socks/protocol.rs) reduces an address to its canonical IP and its port, leaving out IPv6 flow info and scope. A dual-stack socket reports an IPv4 relay’s replies from the IPv4-mapped IPv6 address, and they still count as the relay’s. The app’s bind closure always binds the relay’s own family, so this makes no difference in the app.
  • Both directions first drain the control stream into a 256-byte sink. The control stream reaching EOF ends the association, with BrokenPipe socks: the control connection closed from then on; a read error on it is returned as is.
  • The handshake’s own refusals are PermissionDenied (auth method not supported, server rejects account), ConnectionRefused (server rejects request: <status>), Unsupported (socks: the relay address is a domain) or UnexpectedEof (socks: server closed during the handshake). An Etemenanki upstream answers 0x02 when its relay is bound in an address family the control connection’s address is not in, which arrives here as server rejects request: 2.

The SOCKS UDP link and the client runtime are covered in SOCKS.

app/src/outbound/freedom.rs
#[derive(Clone)]
pub struct FreedomConnector {
dialer: Dialer,
resolver: Resolver,
strategy: AddressFamilyStrategy,
}
impl FreedomConnector {
pub fn new(resolver: Resolver, strategy: AddressFamilyStrategy) -> Self;
}
impl Connector<Flow> for FreedomConnector {
type Stream = TcpStream;
type Datagram = ResolvingUdp;
type Future = DialFuture;
fn connect(&mut self, flow: Flow) -> DialFuture;
}

The connector is cloned into every dial. The clone shares the generation’s Resolver (an Arc inside) and carries the outbound’s address_family strategy. The dialer is Dialer::new(SocketOptions::default()).

TCP. destination_to_socketaddrs(&dest, strategy, &resolver) resolves the destination and orders the candidates by the strategy. dialer.tcp.connect_any(&addrs) then tries them in order, one at a time. Each attempt is bounded by DEFAULT_CONNECT_TIMEOUT (10 s). If every attempt fails, the error is ConnectionRefused failed to connect to any address (…), listing each address and its own error.

UDP. dialer.udp.bind_dual(families) binds up to one socket per family. The families are [V4] for ipv4_only, [V6] for ipv6_only, and both otherwise. A family whose bind fails is left out. Only when nothing binds does the dial fail (AddrNotAvailable, udp: no usable local socket in any requested family). The resulting DualStackUdp is wrapped in a ResolvingUdp:

app/src/outbound/freedom.rs
pub struct ResolvingUdp {
socket: DualStackUdp,
resolver: Resolver,
strategy: AddressFamilyStrategy,
resolved: HashMap<CompactString, Option<IpAddr>>,
resolving: Option<(CompactString, ResolveFuture)>,
}

poll_send_to calls poll_target first:

  • An IP destination goes out at once.
  • A domain already in resolved goes out to the cached address. A cached None means the name did not resolve, and the packet is dropped: poll_send_to returns Ok(buf.len()) and logs freedom: dropping a datagram to an unresolvable … at debug level.
  • A new domain starts one lookup, stored in resolving, and the packet waits (Pending). resolving holds one lookup at a time; a packet to another name returns Pending while it is occupied. The server runtime retries the same packet at the head of its queue until it is sent, so in practice the lookup in flight is the one for that packet, and ordering is preserved.
  • The lookup keeps the first address in the strategy’s order (destination_to_socketaddrs does not know which families bind_dual bound), and the result, success or failure, is cached for the life of the link.

DualStackUdp::poll_send_to then sends from the socket of that address’s family. If that family was not bound, the send fails with AddrNotAvailable (udp: no local socket in the family of <peer>), and the runtime reports it to the core as SendFailed. Because the address is cached, every later packet to that name fails the same way. poll_recv_from reads from whichever socket has a datagram, checking v4 before v6, and returns the source as Destination::udp.

app/src/outbound/udp_fanout.rs
pub const MAX_SUBS: usize = 64;
type OpenFuture = Pin<Box<dyn Future<Output = io::Result<OutboundDatagram>> + Send>>;
struct Sub {
key: usize,
link: OutboundDatagram,
}
pub struct FanOutLink {
router: Arc<Router>,
ctx: FlowContext,
flow: Flow,
subs: Vec<Sub>,
next: usize,
opening: Option<(usize, OpenFuture)>,
recv_waker: Option<Waker>,
}
impl FanOutLink {
pub fn new(router: Arc<Router>, ctx: FlowContext, flow: Flow) -> Self;
}
impl DatagramLink for FanOutLink {
type Addr = Destination;
// poll_send_to, poll_recv_from
}

Every datagram carries its own address, so a UDP association has no single destination to route on. Routing it once, as a unit, would give one outbound every peer’s traffic, and UDP would escape every route rule after the first packet. FanOutLink instead routes each packet and keeps a small table of per-outbound sub-links. It is covered in depth below.

app/src/balancer.rs
pub const DEFAULT_PROBE_INTERVAL: Duration = Duration::from_secs(30);
pub const DEFAULT_PROBE_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Strategy {
Failover,
RoundRobin,
}
impl Strategy {
pub fn parse(s: &str) -> io::Result<Self>;
}
pub struct Member {
pub tag: CompactString,
pub outbound: Arc<Outbound>,
pub probe: Destination,
healthy: AtomicBool,
}
impl Member {
pub fn new(tag: CompactString, outbound: Arc<Outbound>, probe: Destination) -> Self;
pub fn is_healthy(&self) -> bool;
pub fn set_healthy(&self, up: bool) -> bool;
}
pub struct Balancer {
members: Vec<Arc<Member>>,
strategy: Strategy,
next: AtomicUsize,
}
impl Balancer {
pub fn new(members: Vec<Arc<Member>>, strategy: Strategy) -> io::Result<Self>;
pub fn select(&self) -> Arc<Outbound>;
pub fn spawn_probe(
self: &Arc<Self>,
token: CancellationToken,
interval: Duration,
timeout: Duration,
resolver: Resolver,
);
}

A member’s outbound is the same Arc<Outbound> the tag maps to, so a member stays usable on its own under its own tag. Balancing is covered in depth below.

sequenceDiagram
  participant C as Server core
  participant RT as Server runtime
  participant A as AppConnector
  participant R as Router
  participant O as Outbound
  C->>RT: Effect Open with a TCP flow
  RT->>A: connect(flow)
  A->>R: route(route_target)
  R-->>A: Arc of the matched Outbound
  A->>O: connect_stream(flow)
  Note over O: Balanced calls select first
  O-->>A: boxed StreamFuture
  A-->>RT: ConnectFuture
  RT->>RT: poll the future inside the connection task
  alt dial and handshake succeed
    RT->>C: Event Connected
  else any error
    RT->>C: Event ConnectFailed
  end

From here on the runtime polls the OutboundStream directly. For a proxy client, that means polling a ProxyClientRuntime inside the connection’s own task. No task or channel sits between the inbound and the upstream.

sequenceDiagram
  participant C as Server core
  participant RT as Server runtime
  participant F as FanOutLink
  participant R as Router
  participant O as Routed Outbound
  participant S as Sub-link
  C->>RT: Effect Open with a UDP flow
  RT->>F: AppConnector builds the link, no dial
  RT->>C: Event Connected
  C->>RT: Effect SendTo peer
  RT->>F: poll_send_to(packet, peer)
  F->>R: route(flow toward peer)
  R-->>F: Arc of an Outbound, key is its pointer
  alt a sub-link has this key
    F->>S: poll_send_to
  else no sub-link yet
    F->>O: connect_datagram(flow toward peer)
    F-->>RT: Pending, the packet stays queued
    O-->>F: OutboundDatagram, pushed as a new sub-link
    F->>S: poll_send_to
  end
  RT->>F: poll_recv_from
  F->>S: poll each sub-link, starting at next
  S-->>F: payload and source
  F-->>RT: source
  RT->>C: Event Datagram

poll_send_to(cx, buf, to) starts every packet the same way:

let outbound = self
.router
.route(route_target(&self.flow.toward(to.clone()), &self.ctx));
let key = Self::key_of(&outbound);
  • Flow::toward copies the association’s user and source into a flow aimed at this packet’s destination, and clears sniffed. Per-packet routing therefore sees the packet’s own address (an IP or a domain, whatever the client sent), the inbound tag, the source, and the UDP network. A domain or geosite rule matches a UDP packet only when the client addressed that packet by name.
  • The route lookup runs for every packet, not once per peer. Its cost is the first-match walk of the route table described in Routing.

A sub-link is identified by Arc::as_ptr(outbound) as usize. The pointer is a stable name for “the same outbound” because the Router owns every Arc<Outbound> for the whole generation, and the FanOutLink holds an Arc<Router>. No outbound can be freed, and its address reused, while the link exists. Consequences:

  • Two peers routed to the same outbound share one sub-link. For freedom, SOCKS, Trojan, WireGuard and Hysteria 2 that is one socket or one proxied association carrying many peers.
  • A balancer is one key. The member is chosen once, by select() inside connect_datagram, when the sub-link opens. Packets keep using that member’s link until the sub-link is evicted or ends, even if the member’s health changes.
  • A member used both directly (by its own tag) and through a balancer is two keys, so it can get two sub-links.
  • blackhole is a sub-link too: a BlackholeLink that swallows packets. That is how a blocked peer’s packets are dropped without harming the rest of the association.

The sub-link’s datagram link decides how the per-packet address reaches the far end:

Outbound Per-packet destination honoured by the sub-link
freedom Yes: ResolvingUdp sends each packet to its own address.
blackhole Not applicable: packets are swallowed.
socks Yes: every packet carries its own SOCKS5 UDP header.
trojan Yes: every Trojan UDP packet carries its address.
vless, vmess No: the request header names one target, the destination of the packet that opened the sub-link, and seal_to ignores the per-packet address. open_from attributes every reply to that same target.
wireguard, hysteria2 Yes: each datagram is sent to its own address.
http, shadowsocks (both families) No UDP: the open fails and the packet is dropped.
flowchart TB
  A["route the packet, key = outbound pointer"] --> B{"sub-link with this key?"}
  B -- yes --> C["sub.poll_send_to"]
  C --> D{"Ready Ok?"}
  D -- yes --> E["move the sub-link to the back"]
  D -- "no" --> R["return the result as is"]
  E --> R
  B -- no --> F{"opening slot"}
  F -- empty --> G["opening = connect_datagram, then loop"]
  G --> B
  F -- "same key" --> H["poll_opening"]
  H -- Pending --> P["return Pending"]
  H -- "opened" --> B
  H -- failed --> X["drop the packet, return Ok"]
  F -- "another key" --> I["poll_opening"]
  I -- Pending --> P
  I -- Ready --> B

The rules this loop implements:

  • Least recently used order. subs is ordered by the last successful send, most recent at the back. A newly opened sub-link joins at the back before any send. A send that returns anything but Ready(Ok) does not move the sub-link.
  • One open at a time. opening holds at most one (key, OpenFuture). A packet for any other outbound that needs a new sub-link waits for the open in flight: it returns Pending. The server runtime applies effects strictly in order, so the packet stays at the head of the queue, and the association’s uplink waits with it. When the open resolves, the loop runs again and the waiting packet starts its own open.
  • A failed open drops the packet. If the open for this packet’s key fails, poll_send_to returns Ready(Ok(buf.len())). The association does not stall on a peer it cannot reach, and the core does not see an error. The failure is logged at debug level as udp fan-out: opening an outbound failed: <error>. Nothing is cached: the next packet routed to that outbound starts a new open. poll_recv_from also drives the open (see below); when it is the receive side that sees the failure, it only logs it and clears the slot, so the waiting packet finds no open in flight on its next attempt and starts a fresh one instead of being dropped. This is also what happens to UDP routed to http or shadowsocks, whose connect_datagram fails with Unsupported.
  • A send error is reported, not fatal. An error from a sub-link’s poll_send_to is returned as is. The server runtime turns it into Event::SendFailed for that packet, and the fan-out stays open.
stateDiagram-v2
  [*] --> Idle
  Idle --> Opening: a packet routes to an outbound with no sub-link
  Opening --> Opening: polled by a send or a receive, still pending
  Opening --> Idle: open succeeded, sub-link pushed, receiver woken
  Opening --> Idle: open failed, logged

poll_opening is the only place a sub-link joins the table. On success it:

  1. pushes the new Sub to the back of subs;
  2. if subs.len() > MAX_SUBS, removes index 0, the sub-link least recently sent to (or opened), and decrements next if it was past 0. Dropping the removed OutboundDatagram closes its socket, proxy stream or association;
  3. wakes recv_waker, so a receiver parked on an empty table polls the new sub-link.

The table therefore never holds more than MAX_SUBS (64) open sub-links, plus the one being opened.

poll_recv_from first calls poll_opening, so an open makes progress even while no packet is being sent. Then it polls every sub-link once, starting at next and wrapping around:

  • Ready(Ok(from)): sets next to the following index and returns. The next receive starts after this sub-link, so one busy peer cannot starve the others.
  • Ready(Err(e)): the sub-link has ended (for example, a SOCKS control stream closed or a proxy stream failed). It is removed, next is reset to 0 if it fell off the end, and the task wakes itself so the remaining sub-links are polled again. The error is logged at debug level as udp fan-out: a sub-link ended: <error> and never returned to the runtime.
  • Pending from all of them: stores cx.waker() in recv_waker and returns Pending.

FanOutLink::poll_recv_from never returns an error. The server runtime treats a receive error as fatal for the outbound, and here the outbound is the whole association. One failing peer must not end it.

build in app/src/instance.rs checks each [[balancer]] before constructing it. All of these fail the config, so --test catches them:

Check Error
The balancer tag is not already in the outbound map (outbounds and earlier balancers) balancer tag <tag> collides with an outbound tag
Every member names an [[outbound]] balancer <tag> references unknown outbound tag: <member>
Every member has an upstream a TCP probe can reach (upstream_dest_opt) balancer <tag>: outbound <member> has no upstream a TCP health probe can reach, so it cannot be balanced
strategy, when set, is exactly failover or round_robin (unset means failover) unknown balancer strategy "<value>" (expected "failover" or "round_robin")
outbounds is not empty a balancer needs at least one outbound

upstream_dest_opt returns the outbound’s server and port as a TCP Destination. It returns None when either is missing. freedom, blackhole and wireguard do not use those two keys, so they are refused as long as the keys are left out. The check looks only at the keys, though, and build_outbound does not reject server and port on these protocols: with both set, --test accepts such a member, and its probe dials an address the outbound never uses. It also returns None for hysteria2 (and its aliases hysteria and hy2): that server listens on UDP only, so a TCP probe would mark it down on the first cycle and keep it down.

The second check looks the member up in the config’s [[outbound]] list, not in the outbound map. A balancer therefore cannot be a member of another balancer, and the error for that case is the “unknown outbound tag” one. A member may be listed twice; it then gets two Member entries, each with its own probe task.

spawn_generation calls spawn_probe for every balancer before it binds any inbound, with the generation’s token and the balancer’s probe_interval and probe_timeout (defaults DEFAULT_PROBE_INTERVAL, 30 s, and DEFAULT_PROBE_TIMEOUT, 5 s). spawn_probe starts one detached tokio::spawn task per member:

stateDiagram-v2
  [*] --> Probing
  Probing --> Recording: connect succeeded or failed, bounded by timeout
  Recording --> Sleeping: set_healthy, log if the state changed
  Sleeping --> Probing: interval elapsed
  Sleeping --> [*]: generation token cancelled
  • probe resolves member.probe with AddressFamilyStrategy::Auto through the generation’s resolver, then tries TcpStream::connect on each address in turn. It is up if any address accepts. tokio::time::timeout bounds the whole thing, resolution included, by timeout. The probe does not use the outbound’s transport, credentials or address_family: it answers “is the upstream reachable”, not “does the proxy work”.
  • set_healthy swaps the AtomicBool (Ordering::Relaxed) and returns the old value. A change is logged at info level as balancer member <tag> is now up or … is now down.
  • The first probe runs as soon as the task starts. Until it finishes, every member counts as healthy, because Member::new starts at true. Starting members as down would drop all traffic for a whole interval after every start or reload.
  • The task checks the token only while sleeping. After a cancel, a task that is mid-probe finishes that probe (at most timeout) and then exits at the select!.

select() runs once per TCP flow, from connect_stream, and once per fan-out sub-link, from connect_datagram. It never runs at route time, so a member going down affects the next flow, not the next reload.

flowchart TB
  A["collect members whose healthy flag is set"] --> B{"any healthy?"}
  B -- no --> F["first configured member"]
  B -- yes --> C{"strategy"}
  C -- Failover --> D["first healthy member in configured order"]
  C -- RoundRobin --> E["healthy at next.fetch_add(1) mod healthy count"]
Situation failover round_robin
Several members healthy The first healthy one in configured order. When a higher member recovers, new flows go back to it. Each healthy member in turn.
Some members down Skipped. Skipped. The counter keeps running, so the rotation shifts when the healthy set changes.
Every member down The first configured member. The first configured member.

Falling back to the first member rather than failing the flow is deliberate. The upstream may have recovered since the last probe, and a balancer that refuses every flow after a bad probe cycle is a black hole. Selection is lock-free: one relaxed atomic load per member into a short Vec of the healthy ones, plus one fetch_add under round_robin when any member is healthy.

Balancing never retries. Dispatch consumes the flow, so a dial that fails on the chosen member reaches the core as ConnectFailed (TCP) or a dropped packet (UDP). Routing around a dead upstream depends entirely on the probe having noticed.

Health state belongs to the generation. A reload builds new Balancer and Member objects with every member healthy, cancels the old generation’s token, which stops the old probe tasks at their next sleep, and then spawn_generation starts the new probe tasks. A reload whose config fails to parse or build leaves the old generation, and its probes, running. See Generations and reload.

Invariant Enforced by Pinned by
A UDP association is routed per packet, not as a unit AppConnector::connect returns a FanOutLink for UDP flows; FanOutLink::poll_send_to calls router.route for every packet one_association_routes_each_peer_separately (e2e_udp_route.rs)
A blocked peer does not harm the rest of its association Blocked packets go to a BlackholeLink sub-link; poll_recv_from never returns an error one_association_routes_each_peer_separately (third exchange)
Replies are attributed to the peer that sent them Each sub-link’s poll_recv_from returns the source; FanOutLink passes it through replies_from_several_peers_merge_back_correctly (e2e_udp_route.rs)
UDP rules see the UDP network route_target sets TargetNetwork::Udp from the packet’s destination network_separates_tcp_from_udp (e2e_route_context.rs)
At most MAX_SUBS sub-links are open per association poll_opening evicts index 0 after a push that exceeds MAX_SUBS not tested directly
At most one sub-link opens at a time opening is a single Option<(usize, OpenFuture)>; other keys wait on it not tested directly
A failed open drops at most its packet poll_send_to returns Ready(Ok(buf.len())) when it sees the open for its key fail not tested directly
The sub-link key stays unique for the link’s life Keys are Arc<Outbound> pointers, and the FanOutLink holds the Arc<Router> that owns them by construction
The ProxyClient lock is never held across I/O ProxyClient::connect returns the dial future; it is awaited outside the lock by construction
A proxy dial resolves only when the flow is really open ProxyClientConnecting waits for dial, handshake and flush server_runtime_relays_through_a_client_runtime_in_one_task (concepts/tests/client.rs)
SOCKS UDP replies are accepted only from the relay SocksUdpLink::poll_recv_from skips a datagram when endpoint(from) != endpoint(self.relay) udp_link_ignores_datagrams_not_from_the_relay, udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay (protocols/tests/pipeline/socks.rs)
Outbounds that cannot carry UDP say so instead of misbehaving connect_datagram returns Unsupported for http, shadowsocks and shadowsocks-2022 not tested directly
Only outbounds with a TCP probe address can be balanced upstream_dest_opt is None without server and port, and for Hysteria 2 hysteria2_cannot_be_balanced (app/tests/unit/outbound.rs), a_member_with_no_upstream_is_refused (e2e_balancer.rs)
failover takes the first healthy member, and returns to a recovered one select with Strategy::Failover failover_takes_the_first_healthy_member_in_order (app/tests/unit/balancer.rs)
round_robin skips members that are down select rotates over the healthy subset round_robin_cycles_only_through_healthy_members
A balancer with every member down still selects select falls back to members.first() every_member_down_still_selects_rather_than_dropping
A balancer is never empty Balancer::new refuses an empty members a_balancer_needs_at_least_one_member
Strategy names are exact Strategy::parse accepts only failover and round_robin strategy_names_are_validated
A member that stops answering stops receiving flows The probe task marks it down; select skips it traffic_moves_off_a_member_that_stops_answering (e2e_balancer.rs)
Probes die with their generation Each probe task selects on the generation’s CancellationToken not tested directly
Where Condition Kind Message
connect_datagram http, shadowsocks, shadowsocks-2022 Unsupported http carries no datagrams (and the same with the protocol name)
proxy_stream, proxy_datagram The client produced the other kind Other a TCP flow was dialed as datagrams, a UDP flow was dialed as a stream
connect_stream, connect_datagram freedom, wireguard or hysteria2 produced the other kind Other freedom dialed a TCP flow as UDP and the matching variants
FreedomConnector, TCP One address does not answer in time TimedOut connect to <addr> timed out (collected, not returned alone)
FreedomConnector, TCP No address connects ConnectionRefused failed to connect to any address (<addr>: <error>; …)
FreedomConnector, UDP No family binds AddrNotAvailable udp: no usable local socket in any requested family
ResolvingUdp send The address’s family has no socket AddrNotAvailable udp: no local socket in the family of <peer>
ResolvingUdp send The name does not resolve none Packet dropped; debug log freedom: dropping a datagram to an unresolvable …
SocksOutbound UDP Control stream closed BrokenPipe socks: the control connection closed
FanOutLink send The open for this packet’s outbound fails none Packet dropped; debug log udp fan-out: opening an outbound failed: <error>
FanOutLink receive A sub-link’s receive fails none Sub-link removed; debug log udp fan-out: a sub-link ended: <error>

How the server runtime treats each outcome:

  • An Err from a connect_stream future becomes Event::ConnectFailed.
  • An Err from a sub-link send becomes Event::SendFailed, and the association stays open.
  • Because FanOutLink swallows open failures and sub-link receive errors, a UDP association ends only when its inbound side ends it.

Nothing in this layer owns a task except the balancer probes. Cancellation is by drop:

  • Dropping a StreamFuture or DatagramFuture before it resolves drops the transport dial, the handshake state and any partly built stream.
  • Dropping a FanOutLink drops every sub-link and the open in flight, if any. This happens when the connection’s runtime ends or the generation is cancelled. Each sub-link closes its own sockets or streams as it drops.
  • Dropping an OutboundStream closes the TCP stream, or drops the client runtime and with it the upstream.
  • Probe tasks are detached. No JoinHandle is kept. They end on the generation token, as described above.
Constant Where Value Meaning
MAX_SUBS app/src/outbound/udp_fanout.rs 64 Open sub-links per UDP association; the least recently sent-to one is closed first
DEFAULT_PROBE_INTERVAL app/src/balancer.rs 30 s Sleep between probes of one member, when probe_interval is not set
DEFAULT_PROBE_TIMEOUT app/src/balancer.rs 5 s Bound on one probe (resolution and every connect attempt), when probe_timeout is not set
DEFAULT_CONNECT_TIMEOUT environment/src/dial/tcp.rs 10 s One TCP connect attempt by TcpDialer: freedom TCP and every TransportConnector dial
RECV_BUF protocols/src/socks/udp_link.rs 64 * 1024 Largest SOCKS relay packet read from the UDP socket
HTTP_BUF, SOCKS_BUF, TROJAN_BUF, VLESS_BUF app/src/outbound/mod.rs 16 * 1024 Client runtime buffer size for these protocols
SS_BUF app/src/outbound/mod.rs 20 * 1024 Shadowsocks (legacy AEAD) client runtime
VMESS_BUF, SS2022_BUF app/src/outbound/mod.rs 32 * 1024 VMess and Shadowsocks 2022 client runtimes

The configured probe_interval and probe_timeout are whole seconds (u64). build passes them to Duration::from_secs without a range check, so --test accepts 0 for either.

Per UDP association, MAX_SUBS bounds the number of open sub-links. Each sub-link costs what its link type costs: two client buffers for a proxy client, a 64 KiB receive buffer for SOCKS, one or two sockets for freedom.

Unit tests live in the binary crate (#[path] modules under app/tests/unit/). The end-to-end tests are one integration crate, app/tests/integration.rs, that spawns the real binary. The SocksUdpLink tests are in the etemenanki-protocols pipeline crate, protocols/tests/pipeline.rs.

Terminal window
cargo test -p etemenanki-app --bin etemenanki-app balancer::
cargo test -p etemenanki-app --bin etemenanki-app outbound::
cargo test -p etemenanki-app --test integration e2e_udp_route
cargo test -p etemenanki-app --test integration e2e_balancer
cargo test -p etemenanki-protocols --test pipeline pipeline::socks::
Test File Behaviour pinned
failover_takes_the_first_healthy_member_in_order app/tests/unit/balancer.rs With a down, b is picked. After a recovers, it is picked again.
round_robin_cycles_only_through_healthy_members app/tests/unit/balancer.rs With b down, four picks give a, c, a, c.
every_member_down_still_selects_rather_than_dropping app/tests/unit/balancer.rs All down: the first member is still returned.
a_balancer_needs_at_least_one_member app/tests/unit/balancer.rs Balancer::new refuses an empty list.
strategy_names_are_validated app/tests/unit/balancer.rs failover and round_robin parse; roundrobin does not.
hysteria2_cannot_be_balanced app/tests/unit/outbound.rs upstream_dest_opt is None for hysteria2, hysteria and hy2, and Some for a Trojan outbound.
hysteria2_refuses_a_stream_block and the other hysteria2_* tests app/tests/unit/outbound.rs build_outbound accepts a minimal Hysteria 2 config and its aliases, and fails closed on bad Hysteria 2 settings.
wireguard_ipv4_only_builds_with_ipv4_address, wireguard_ipv6_only_requires_ipv6_address app/tests/unit/outbound.rs address_family must match the WireGuard interface addresses.
one_association_routes_each_peer_separately app/tests/integration/e2e_udp_route.rs One SOCKS association, three packets: the allowed peer answers, the port-blocked peer does not, and the allowed peer still answers afterwards.
replies_from_several_peers_merge_back_correctly app/tests/integration/e2e_udp_route.rs Two peers on one association each get their own reply, attributed to the right source.
network_separates_tcp_from_udp app/tests/integration/e2e_route_context.rs A network = "udp" rule blocks UDP packets and leaves TCP alone.
traffic_moves_off_a_member_that_stops_answering app/tests/integration/e2e_balancer.rs With probe_interval = 1, a first member that accepts TCP but is not a proxy is used (the flow fails). Once it stops accepting, the probe marks it down and flows move to the working member.
a_member_with_no_upstream_is_refused app/tests/integration/e2e_balancer.rs --test rejects a balancer over a freedom outbound.
new_server_vs_new_client_udp protocols/tests/pipeline/socks.rs SocksUdpLink round-trips datagrams through the SOCKS server core.
udp_link_ignores_datagrams_not_from_the_relay protocols/tests/pipeline/socks.rs Against a hand-written server, a well-formed reply sent to the client’s socket from a socket other than the relay is dropped; the relay’s own reply that follows is returned, attributed to the peer its header names.
udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay protocols/tests/pipeline/socks.rs Linux only. A link whose bind closure binds [::]:0 associates with an IPv4 relay and still receives its replies, which arrive from the IPv4-mapped address.

None of the end-to-end tests above need a Go toolchain. The app’s own SOCKS inbound produces the association, and a second app instance plays the upstream. The app_client_* tests in app/tests/integration/e2e_xray.rs run the VLESS client against a real Xray server over WebSocket and gRPC (each plain and with TLS) and over TLS on TCP, and skip themselves when go is not installed.

Several fan-out properties (the MAX_SUBS eviction, the single open slot, the drop on a failed open) have no dedicated test. If you change FanOutLink, a unit test with a stub Router and scripted OutboundDatagrams is the missing piece.