Outbounds, UDP fan-out and balancers
Source files: 72 · checked against Etemenanki 555b7df
Etemenanki/supervisor/src/build/outbound.rsEtemenanki/supervisor/src/topology/outbound/mod.rsEtemenanki/supervisor/src/topology/outbound/proxy.rsEtemenanki/supervisor/src/topology/outbound/freedom.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/topology/spec_plan/outbound.rsEtemenanki/supervisor/src/topology/spec_plan/transport.rsEtemenanki/supervisor/src/topology/spec_plan/route.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/metered.rsEtemenanki/concepts/src/client.rsEtemenanki/concepts/src/link.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/socks/udp_link.rsEtemenanki/protocols/src/socks/codec.rsEtemenanki/protocols/src/http/codec.rsEtemenanki/protocols/src/trojan/codec.rsEtemenanki/protocols/src/trojan/protocol.rsEtemenanki/protocols/src/vless/codec.rsEtemenanki/protocols/src/vless/protocol.rsEtemenanki/protocols/src/vmess/codec.rsEtemenanki/protocols/src/ss_legacy/codec.rsEtemenanki/protocols/src/ss_2022/codec.rsEtemenanki/protocols/src/ss_2022/crypto.rsEtemenanki/protocols/src/ss_aead.rsEtemenanki/protocols/src/helpers/address.rsEtemenanki/protocols/src/helpers/address_family.rsEtemenanki/protocols/src/transports/connect.rsEtemenanki/protocols/src/transports/keepalive.rsEtemenanki/protocols/src/transports/tls/config.rsEtemenanki/protocols/src/hysteria/config.rsEtemenanki/protocols/src/hysteria/connector.rsEtemenanki/protocols/src/hysteria/connection.rsEtemenanki/protocols/src/hysteria/slot.rsEtemenanki/protocols/src/wireguard/config.rsEtemenanki/protocols/src/wireguard/connector.rsEtemenanki/protocols/src/wireguard/slot.rsEtemenanki/protocols/src/wireguard/device.rsEtemenanki/environment/src/dial/mod.rsEtemenanki/environment/src/dial/tcp.rsEtemenanki/environment/src/dial/udp.rsEtemenanki/environment/src/dial/socket.rsEtemenanki/app/src/lower.rsEtemenanki/app/src/subscribe.rsEtemenanki/ffi/src/platform.rsEtemenanki/supervisor/tests/unit/balancer.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/validate.rsEtemenanki/supervisor/tests/socket_policy.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/concepts/tests/client.rsEtemenanki/protocols/tests/pipeline/socks.rsEtemenanki/app/tests/integration/e2e_udp_route.rsEtemenanki/app/tests/integration/e2e_balancer.rsEtemenanki/app/tests/integration/e2e_route_context.rsEtemenanki/app/tests/integration/e2e_xray.rsEtemenanki/app/tests/integration/e2e_xray_vmess.rsEtemenanki/app/tests/integration/e2e_hysteria.rsEtemenanki/app/tests/integration/e2e_wg.rs
An outbound is where a routed flow leaves the process. The supervisor builds one handler per OutboundSpec when an apply needs one, keeps it for as long as its spec and the DNS spec stay the same, and opens flows on it through two methods: one for a TCP flow, one for a UDP flow. This page covers that layer in etemenanki-supervisor: build_outbound and the resolvers and socket policy each protocol is given, the closed Outbound enum and the stream and datagram types it hands back, the proxy clients and their buffers, the SOCKS, freedom, blackhole, Hysteria 2 and WireGuard outbounds, the per-packet UDP fan-out, balancers, and the socket policy on every socket an outbound opens.
It is written for contributors who add an outbound protocol, change how a UDP association is routed, or touch balancing. How a flow reaches a target (Plane, Target, AppConnector::connect) and what a drain does to open flows are on the routing plane. The client runtime every proxy outbound runs on is on Client codecs and the client runtime. The operator’s view of the same features is in Outbounds and Balancers.
Responsibilities
Section titled “Responsibilities”| Component | File → symbol | Owns |
|---|---|---|
| Construction | supervisor/src/build/outbound.rs → build_outbound, transport |
Turning one OutboundSpec into one Outbound: the per-flow codec closures, the transport connector, the resolvers and the socket policy. |
| Dispatch | supervisor/src/topology/outbound/mod.rs → Outbound |
Opening a TCP flow as a byte stream (connect_stream) and a UDP flow as a datagram link (connect_datagram), and refusing UDP where the protocol has none. |
| Proxy clients | supervisor/src/topology/outbound/proxy.rs → ProxyClient |
One ProxyClientConnector per outbound behind a mutex, and the closed OutboundStream, OutboundDatagram and BlackholeLink types. |
| SOCKS | supervisor/src/topology/outbound/mod.rs → SocksOutbound |
CONNECT through a ProxyClient, and UDP ASSOCIATE as a link of its own over a control stream. |
| Direct | supervisor/src/topology/outbound/freedom.rs → FreedomConnector, ResolvingUdp |
Resolving and connecting TCP; binding UDP sockets per family and resolving domain destinations. |
| UDP fan-out | supervisor/src/topology/outbound/udp_fanout.rs → FanOutLink |
Routing every packet of a UDP association on the current plane, keeping one sub-link per outbound version, and merging the replies. |
| Balancers | supervisor/src/topology/balancer.rs → Balancer, Member, Strategy |
Health per member, one probe task per member, and choosing a member per flow or per packet. |
What it leaves to others:
- Routing.
Plane::route,Target,GuardedandAppConnector::connectare on the routing plane. This page usesTarget::resolve,Target::connect_datagramandGuardedonly as far as the fan-out and the balancers need them. - Versions, reuse and drains. Which outbounds an apply builds, reuses or drains is decided by the plan; see Reconcile.
- Metering, kills and pacing. Every stream and every fan-out sub-link is a tracked flow; see Tracking.
- The DNS service and the resolvers.
Outbound::Dns,DnsUdpLink,DnsTcpStream,RoutedDialerand theDnsresolver pair are on Supervisor DNS. - Spec rules. Everything
validaterefuses is on Validation; this page lists only the balancer checks. - Wire formats. Each codec and client is on its protocol page, for example SOCKS, Hysteria 2: client and WireGuard.
Where outbounds live
Section titled “Where outbounds live”plan(supervisor/src/topology/spec_plan/plan.rs) gives every outbound tag a version; the pair is anOutboundId { tag, version }(supervisor/src/entity/id.rs). An outbound whose spec is unchanged is aStep::Reuseat its running version. Anything else is aStep::Buildat the next version, and the old version gets aStep::Drain. A changed DNS spec rebuilds every outbound exceptblackhole, because every other outbound holds resolvers built from it (uses_dns). Versions of removed tags are kept inRunningState::versions, so a tag added back never reuses a version whose flows may still be draining.Actor::prepare(supervisor/src/supervisor.rs) callsbuild_outbound(outbound, &dns, &self.socket)for every built outbound and wraps the result withTarget::outbound(id, built)in anArc. A reused outbound is the runningArc<Target>itself, so it keeps what it holds: itsProxyClient, its QUIC connection or its tunnel.- Balancers are built next, because a member is the
Arc<Target>of an outbound built or reused in the same apply (see State across applies). commitpublishes the new plane, then stops the probes of balancers that were rebuilt or removed and starts probes for the new ones.
A construction failure is an ApplyError::Build, printed building outbound <tag>@v<n> failed: <error>, and the whole apply is refused while the running state stays as it was. etemenanki-app prints it after configuration invalid: . For a TLS outbound whose ca_file holds no certificate, the pinned binary prints:
configuration invalid: building outbound up@v1 failed: no certificate in CA PEM bundlecheck (the app’s --test) runs prepare without binding. It builds every outbound and every balancer, so construction errors are caught, but no outbound opens a socket while it is constructed, and no probe starts. The first socket an outbound opens belongs to its first flow.
Key types
Section titled “Key types”build_outbound
Section titled “build_outbound”pub(crate) fn build_outbound( spec: &OutboundSpec, dns: &Dns, socket: &SocketOptions,) -> io::Result<Outbound>;
pub(crate) fn obfs(spec: &ObfsSpec) -> Obfs;fn tcp(upstream: &ProxyUpstream) -> Destination;fn udp(server: &Destination) -> Destination;fn transport( spec: &OutboundTransportSpec, dns: &Dns, socket: &SocketOptions,) -> io::Result<TransportConnector>;Dns (supervisor/src/build/dns.rs) carries two resolvers. servers looks up the proxy servers the outbounds dial. destinations looks up the destinations this process connects to itself: through freedom, and inside a WireGuard tunnel. With DnsSpec::Single both are one resolver. With DnsSpec::Split, servers asks the pre_proxy servers directly and destinations asks the through_proxy servers through the route table; see Supervisor DNS.
socket is the supervisor’s socket policy, handed to every dialer a protocol is built with (see Socket policy on every outbound socket).
What each spec becomes:
OutboundProtocolSpec |
Outbound variant |
Built from |
|---|---|---|
Freedom { address_family } |
Freedom(FreedomConnector) |
FreedomConnector::new(Dialer::new(socket), dns.destinations, address_family) |
Blackhole |
Blackhole |
nothing |
Socks { upstream, account } |
Socks(Box<SocksOutbound>) |
SocksOutbound::new(transport(..), UdpDialer::new(socket), tcp(upstream), auth) |
Http { upstream, account } |
Http(ProxyClient<HTTP_BUF, HttpConnect, NoUdp>) |
make builds HttpConnect::new(&flow.destination, auth) |
Trojan { upstream, password } |
Trojan(ProxyClient<TROJAN_BUF, TrojanStream, TrojanDatagram>) |
password_hash(password) once; make builds TrojanStream::new(&hash, &dest) or TrojanDatagram::new(&hash) |
Vless { upstream, id } |
Vless(ProxyClient<VLESS_BUF, VlessStream, VlessDatagram>) |
make builds VlessStream::new(&uuid, &dest) or VlessDatagram::new(&uuid, &dest) |
Vmess { upstream, id, security } |
Vmess(ProxyClient<VMESS_BUF, VMessStream, VMessDatagram>) |
make builds VMessStream::new(uuid, security, true, &dest) or VMessDatagram::new(..); true is global_padding |
Shadowsocks { upstream, method, password } |
Shadowsocks(ProxyClient<SS_BUF, SsStream, NoUdp>) |
evp_bytes_to_key(password, method.key_len()) once; make builds SsStream::new(method, key.clone(), &dest) |
Ss2022 { upstream, method, identity_psks, psk } |
Ss2022(ProxyClient<SS2022_BUF, Ss2022Stream, NoUdp>) |
make builds Ss2022Stream::new(method, psk.clone(), identity.clone(), &dest) |
Hysteria2(spec) |
Hysteria2(Hy2Connector) |
Hy2Connector::with_address_family(Hy2Config { .. }, address_family).with_resolver(dns.servers) |
Wireguard(spec) |
Wireguard(WgConnector) |
WgConnector::with_address_family(WgConfig { .. }, address_family).with_resolver(dns.destinations) |
Which resolver and which sockets each one uses:
| Protocol | Proxy server or endpoint resolved by | Destinations resolved by | Sockets opened under the policy |
|---|---|---|---|
freedom |
none | destinations, under the outbound’s address_family |
TCP connects; one UDP socket per family |
blackhole |
none | none | none |
socks |
servers, under the transport’s address_family |
the upstream | the TCP connect under the transport; the UDP relay socket |
http, trojan, vless, vmess, shadowsocks (both families) |
servers, under the transport’s address_family |
the upstream | the TCP connect under the transport |
hysteria2 |
servers, under the outbound’s address_family |
the server | the QUIC UDP socket |
wireguard |
servers, under endpoint_address_family |
destinations, under address_family and the tunnel’s own address families |
the tunnel’s UDP socket |
Three details of the table:
- Keys are derived once. The Trojan password is hashed and the Shadowsocks master key derived when the outbound is built, and the
makeclosures capture the result. Only the per-flow codec is built per flow. tcpandudpfix the network.tcp(upstream)copies the proxy server’sDestinationwithnetwork = DialNetwork::Tcp;udp(server)does the same withDialNetwork::Udpfor the Hysteria 2 server and the WireGuard endpoint.obfsconverts the spec’sObfsSpec::Salamander { psk }into the protocols crate’sObfs::Salamander { psk }.
The transport under a proxy client
Section titled “The transport under a proxy client”transport(spec, dns, socket) maps the upstream’s StreamShape to a TransportKind. tls(alpn) is a closure over OutboundTransportSpec::tls that returns an Option<ClientConfig>: Some when the spec carries TLS material, None otherwise.
StreamShape |
TransportKind built |
TLS layered | Error when a field is missing |
|---|---|---|---|
Tcp |
TransportKind::Tcp |
never | none |
Tls |
TransportKind::Tls(config), from tls(Alpn::None) |
always, no ALPN | missing tls |
Ws { path, host, .. } |
TransportKind::ws(host, path, tls(Alpn::Http1)?) |
when OutboundTransportSpec::tls is set, ALPN http/1.1 |
missing ws host |
Grpc { service, authority, .. } |
TransportKind::grpc(authority, service, tls(Alpn::Http2)?) |
when OutboundTransportSpec::tls is set, ALPN h2 |
missing grpc authority |
transport ignores the shape’s own tls flag: TLS is layered exactly when OutboundTransportSpec::tls is set, and validation keeps the two equal (outbound_transport refuses a spec where shape.uses_tls() != tls.is_some()). The three missing … errors are InvalidInput and cannot occur for a validated spec: an outbound’s transport spec carries no gaps, because the front end fills in the WebSocket host, the gRPC authority and the SNI, and validation checks that they are there.
TransportKind::grpc always builds GrpcMode::Gun (one Hunk per gRPC message) and opens with the user-agent DEFAULT_USER_AGENT, a desktop Chrome string (protocols/src/transports/grpc/settings.rs). The protocols crate also offers multi() for MultiHunk framing and user_agent() to change or drop the header, but the supervisor calls neither, so every gRPC outbound uses gun framing and the default user agent.
The TLS configuration is ClientConfig::with_verify_mode(&tls.server_name, tls.verify.clone(), alpn) (protocols/src/transports/tls/config.rs). It builds an OpenSSL SslConnector with TLS 1.2 as the minimum version, and sets verification from VerifyMode:
VerifyMode |
Trust | Hostname checked |
|---|---|---|
System |
the default verify paths (system roots) | yes |
CustomCa(pem) |
the default verify paths plus every certificate of the bundle | yes |
Insecure |
none: SslVerifyMode::NONE |
no |
The CustomCa bundle is parsed here, when the outbound is built. A bundle that does not parse fails with OpenSSL’s own error text (kind Other), and one that parses to no certificate fails with no certificate in CA PEM bundle (InvalidInput).
The result is TransportConnector::new(kind, Dialer::new(socket.clone()), dns.servers.clone(), spec.address_family). TransportConnector::dial (protocols/src/transports/connect.rs) runs these steps:
- It refuses a UDP destination with
Unsupporteda proxy transport carries no datagrams of its own. The proxy clients always hand ittcp(upstream), and a proxy’s UDP rides inside its stream. - It resolves the proxy server with
serversunder the transport’saddress_family(destination_to_socketaddrs). - It connects to the first address that answers (
TcpDialer::connect_any). - It sets TCP keepalive on the connection with
set_keepalive(protocols/src/transports/keepalive.rs): the first probe afterTCP_KEEPALIVE_IDLE(120 s) of silence, then everyTCP_KEEPALIVE_INTERVAL(30 s), and the connection is given up afterTCP_KEEPALIVE_RETRIES(3) unanswered probes. The comment on the constants gives the reason: Linux waits two hours before its first probe, and a proxy’s peers routinely vanish without a FIN, for example after a mobile NAT rebind. A platform that rejects the options is not an error; the failure is logged at debug level ascould not enable TCP keepalive: <error>. - It wraps the stream in the transport.
The transports are covered in TCP and TLS transports and WebSocket and gRPC transports.
The Outbound enum
Section titled “The Outbound enum”const HTTP_BUF: usize = 16 * 1024;const SOCKS_BUF: usize = 16 * 1024;const TROJAN_BUF: usize = 16 * 1024;const VLESS_BUF: usize = 16 * 1024;const VMESS_BUF: usize = 32 * 1024;const SS_BUF: usize = SsStream::BUF_SIZE;const SS2022_BUF: usize = Ss2022Stream::BUF_SIZE;
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), Dns(etemenanki_protocols::dns::serve::Service),}
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;}Do not confuse it with etemenanki_concepts::link::Outbound<S, D>, the two-armed “stream or datagram” result every Connector returns; the code below names that one link::Outbound. Flow is etemenanki_protocols::flow::Flow<Principal> (supervisor/src/topology/flow.rs).
The enum is closed on purpose. A Target owns one Outbound, and the server runtime stores every outbound it opens by value, so each kind of flow needs one concrete type. Both methods take &self, because every flow routed to a target dials the same value at once. Each arm builds its dial future synchronously and boxes it; no socket is touched until the future is first polled.
| Variant | connect_stream |
connect_datagram |
|---|---|---|
Freedom |
Clones the connector and dials; a stream becomes OutboundStream::Tcp |
Clones the connector and dials; the ResolvingUdp link is boxed into OutboundDatagram |
Blackhole |
Ready at once: OutboundStream::Blackhole |
Ready at once: OutboundDatagram over BlackholeLink |
Socks |
SocksOutbound.tcp.connect(flow), then proxy_stream |
SocksOutbound::connect_datagram(), a UDP ASSOCIATE; the flow is not used |
Http |
ProxyClient::connect, then proxy_stream |
Fails: http carries no datagrams |
Trojan, Vless, Vmess |
ProxyClient::connect, then proxy_stream |
ProxyClient::connect, then proxy_datagram |
Shadowsocks, Ss2022 |
ProxyClient::connect, then proxy_stream |
Fails: shadowsocks carries no datagrams, shadowsocks-2022 carries no datagrams |
Wireguard |
Clones the connector and dials; a stream becomes OutboundStream::Wg |
Clones the connector and dials; the link is boxed into OutboundDatagram |
Hysteria2 |
Clones the connector and dials; a stream becomes OutboundStream::Hy2 |
Clones the connector and dials; the link is boxed into OutboundDatagram |
Dns |
DnsTcpStream::new(service) boxed into OutboundStream::Proxy |
DnsUdpLink::new(service) boxed into OutboundDatagram |
The “no datagrams” errors come from no_udp(what) and have kind Unsupported. FreedomConnector, WgConnector and Hy2Connector are cloned per dial because their connect takes &mut self; the clones are cheap. A WgConnector clone shares its tunnel slot (Arc<tokio::sync::Mutex<DeviceSlot>>) and a Hy2Connector clone shares its connection slot (Arc<parking_lot::Mutex<ConnSlot>>), so cloning never opens a second tunnel or connection.
ProxyClient, Make and NoUdp
Section titled “ProxyClient, Make and NoUdp”Seven protocols are clients over a TransportConnector: HTTP, SOCKS CONNECT, Trojan, VLESS, VMess, Shadowsocks and Shadowsocks 2022. All of them go through one wrapper:
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 is a Connector whose connect takes &mut self too. The outbound is shared, so ProxyClient keeps the connector in a parking_lot::Mutex. connect locks, calls the inner connect, and unlocks when it returns. Under the lock the connector runs make, builds the ProxyClientRuntime (the codec plus two boxed buffers of BUF bytes) and asks the TransportConnector for its dial future, which clones the connector into a boxed async block. No I/O and no .await happen under the lock, so flows through one outbound serialise only on building a future, never on the network.
The codec is picked per flow. Each make closure returns link::Outbound::Stream(codec) or link::Outbound::Datagram(codec). Trojan, VLESS and VMess choose from flow.destination.network == DialNetwork::Udp. HTTP, SOCKS, Shadowsocks and Shadowsocks 2022 always return a stream, and NoUdp stands in for their datagram codec. NoCodec (concepts/src/core.rs) is an uninhabited enum, Never(Infallible, PhantomData<..>), so no value of it can exist and its codec methods are never called. A NoUdp protocol never reaches make with a UDP flow, because connect_datagram refuses it first; SOCKS UDP does not use ProxyClient at all.
The dial future. ProxyClientConnecting resolves only once the upstream is dialed, the codec’s handshake is done and the handshake bytes are flushed. A refused upstream is therefore an Err from the future, which the server core sees as ConnectFailed, not a stream that fails later. Two adapters then fold the result into the closed types:
pub fn proxy_stream<const BUF: usize, S, D>( opened: link::Outbound< ProxyClientRuntime<BUF, S, TransportConnector>, ProxyClientRuntime<BUF, D, TransportConnector>, >,) -> io::Result<OutboundStream>;
pub fn proxy_datagram<const BUF: usize, S, D>( opened: link::Outbound< ProxyClientRuntime<BUF, S, TransportConnector>, ProxyClientRuntime<BUF, D, TransportConnector>, >,) -> io::Result<OutboundDatagram>;A runtime of the wrong kind is a TCP flow was dialed as datagrams or a UDP flow was dialed as a stream (kind Other). Neither happens in practice, because make decides the kind from the same network field that chose between connect_stream and connect_datagram.
Client runtime buffers
Section titled “Client runtime buffers”Each BUF constant sizes both buffers of a ProxyClientRuntime: the staging buffer toward the upstream and the read buffer from it. The comment on the constants gives the rule: the largest wire frame each protocol opens, plus what its codec reserves.
| Outbound | Constant | BUF |
Codec | Codec’s STAGING_RESERVE |
Buffers per flow |
|---|---|---|---|---|---|
http |
HTTP_BUF |
16,384 | HttpConnect |
REQUEST_MAX = 1,024 |
32 KiB |
socks (CONNECT) |
SOCKS_BUF |
16,384 | SocksConnect |
528 | 32 KiB |
trojan |
TROJAN_BUF |
16,384 | TrojanStream, TrojanDatagram |
REQUEST_HEADER_MAX = 56 + 2 + 1 + 259 + 2 = 320 |
32 KiB |
vless |
VLESS_BUF |
16,384 | VlessStream, VlessDatagram |
REQUEST_HEADER_MAX = 1 + 16 + 1 + 1 + 259 = 278 |
32 KiB |
vmess |
VMESS_BUF |
32,768 | VMessStream, VMessDatagram |
HEADER_MAX (358) rounded up to 64 = 384 |
64 KiB |
shadowsocks (legacy AEAD) |
SS_BUF = SsStream::BUF_SIZE |
20,480 | SsStream |
32 + 34 + 259 + 61 = 386 | 40 KiB |
shadowsocks (2022- methods) |
SS2022_BUF = Ss2022Stream::BUF_SIZE |
65,569 | Ss2022Stream |
2,048 | 131,138 bytes |
259 is AddressCodec::MAX_LEN (type, length, a 255-byte name, port), and 34 is the AEAD RECORD_OVERHEAD (a sealed two-byte length with its 16-byte tag, plus the payload’s 16-byte tag).
The two Shadowsocks sizes come from their codecs:
SsStream::BUF_SIZEholds the salt and one chunk ofMAX_PAYLOAD(0x3FFFbytes) with its overhead, with room to spare.Ss2022Stream::BUF_SIZEisMAX_RECORD_LEN=MAX_PACKET_SIZE(0xFFFF) +RECORD_OVERHEAD= 65,569. ss-rust and sing seal a whole read of up to0xFFFFbytes into one record, and every other frame is smaller, so a buffer this size always holds the next frame whole.
ProxyClientRuntime::new asserts BUF_SIZE > STAGING_RESERVE (BUF_SIZE must exceed the codec's STAGING_RESERVE or nothing can be sealed); every pair above satisfies it. The buffers bound what one frame may be in each direction:
- From the upstream. A frame that does not fit the read buffer fails the flow with
InvalidDataupstream frame larger than the client runtime's buffer. - Toward the upstream. A stream write is sealed in pieces, so it never overflows. A datagram is sealed whole: a UDP packet whose payload plus the codec’s
STAGING_RESERVEexceedsBUF(more than 16,064 bytes of payload for Trojan, for example) fails its send withInvalidInputframe larger than the client runtime's buffer. In the fan-out, a failed send drops that packet and the sub-link (see Sending).
The buffer mechanics and the runtime’s other errors are on Client codecs and the client runtime; the ones a flow can meet are listed under Errors.
OutboundStream, OutboundDatagram and BlackholeLink
Section titled “OutboundStream, OutboundDatagram and BlackholeLink”pub trait Stream: AsyncRead + AsyncWrite + Send {}impl<T: AsyncRead + AsyncWrite + Send> Stream for T {}
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;OutboundStreamimplementsAsyncReadandAsyncWriteby delegating to its variant. The proxy clients differ in type per protocol and per buffer size, so they are boxed intoProxy; the DNS service’sDnsTcpStreamrides the same variant. Plain TCP, WireGuard and Hysteria 2 streams keep variants of their own and avoid the box.OutboundStream::Blackholereads EOF at once and accepts every write whole;flushandshutdownsucceed at once.OutboundDatagramboxes every datagram link, because the fan-out keeps links of different outbounds side by side.BlackholeLinkreports everypoll_send_toas sent in full and returnsPendingfrom everypoll_recv_from, for ever.
Before a stream or link reaches the server runtime it is wrapped twice: Guarded (it fails once its outbound version is closed by a drain policy, see the routing plane) and then Metered or MeteredDatagram (see Tracking). A TCP flow is a Metered<Guarded<OutboundStream>>; a fan-out sub-link is a MeteredDatagram<Guarded<OutboundDatagram>>.
SocksOutbound
Section titled “SocksOutbound”pub struct SocksOutbound { transport: TransportConnector, /// Binds the socket a `UDP ASSOCIATE` relays through. udp: UdpDialer, server: Destination, auth: Option<(CompactString, CompactString)>, tcp: ProxyClient<SOCKS_BUF, SocksConnect, NoUdp>,}
impl SocksOutbound { pub fn new( transport: TransportConnector, udp: UdpDialer, server: Destination, auth: Option<(CompactString, CompactString)>, ) -> Self;
fn connect_datagram(&self) -> DatagramFuture;}TCP flows go through tcp, an ordinary ProxyClient whose make builds SocksConnect::new(&flow.destination, auth). UDP does not fit the client runtime, because a SOCKS5 association is two connections: a control stream that must stay open, and a UDP socket toward the server’s relay. connect_datagram builds that link itself:
- It dials a control stream with
transport.dial(&server), so the outbound’s transport (TLS, WebSocket, gRPC) and its TCP keepalive apply to the control stream. - It runs
SocksUdpLink::associate(control, auth, bind): method negotiation, the username and password round whenauthis set, thenUDP ASSOCIATE. The method request offers exactly one method (encode_method_request): username and password (0x02) whenauthis set, no authentication (0x00) otherwise. Each round reads the control stream in 512-byte chunks into one buffer until its parser (parse_method_reply,parse_userpass_reply,parse_reply) has a whole reply, and a parser error, such as a wrong version byte, is returned as it is. - The
bindclosure isudp.bind(AddressFamily::of(relay)): one socket under the socket policy, in the relay’s family, on an ephemeral port.UdpDialermakes an IPv6 socket v6-only.
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>;
pub fn relay(&self) -> SocketAddr;}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 the 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 0.0.0.0 port 0
S-->>O: reply naming the relay
O->>O: bind a UDP socket in the relay's family
O-->>F: OutboundDatagram over SocksUdpLink
F->>O: poll_send_to
O->>S: relay header and payload, to the relay
S-->>O: datagram from the relay
O-->>F: payload and the peer its header names
The 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 once the reply names the relay, and a local address would be the wrong one behind NAT. A server that checks sources holds the association to the address the control connection came from and the port of the first datagram, so the datagrams must leave from the IP the server sees on the control connection.
Once associated, the link behaves as follows:
- Every
poll_send_towraps the payload in the SOCKS5 UDP header for that packet’s own destination (encode_udp_packet_into, into a reusedscratchbuffer of initial capacity 2,048), so one association reaches many peers. poll_recv_fromreads into a boxedRECV_BUF(64 KiB) buffer. It skips any datagram whose source is not the relay (endpoint(from) != endpoint(self.relay), which compares the canonical IP and the port, so an IPv4 relay heard on a dual-stack socket as an IPv4-mapped address still counts) and any datagram whose header does not parse. It copies the payload into the caller’s buffer, truncating a payload that does not fit, as UDP does, and returns the peer the header names.- Both directions first drain the control stream into a 256-byte
sink; whatever the server sends on it is discarded. The control stream reaching EOF ends the association: that call and every later one fail withBrokenPipesocks: the control connection closed. A read error on the control stream is returned once as it is, and marks the stream closed, so every later call fails with the sameBrokenPipeerror. - The handshake fails with
PermissionDenied(auth method not supported,server rejects account),ConnectionRefused(server rejects request: <status>),Unsupported(socks: the relay address is a domain) orUnexpectedEof(socks: server closed during the handshake).
The wire format and the server side of the association are on SOCKS.
FreedomConnector and ResolvingUdp
Section titled “FreedomConnector and ResolvingUdp”#[derive(Clone)]pub struct FreedomConnector { dialer: Dialer, resolver: Resolver, strategy: AddressFamilyStrategy,}
impl FreedomConnector { pub fn new(dialer: Dialer, resolver: Resolver, strategy: AddressFamilyStrategy) -> Self;}
impl Default for FreedomConnector; // Dialer::default(), Resolver::default(), Auto
impl Connector<Flow> for FreedomConnector { type Stream = TcpStream; type Datagram = ResolvingUdp; type Future = DialFuture; fn connect(&mut self, flow: Flow) -> DialFuture;}
pub struct ResolvingUdp { socket: DualStackUdp, resolver: Resolver, strategy: AddressFamilyStrategy, // private lookup state}
impl DatagramLink for ResolvingUdp { type Addr = Destination; // poll_send_to, poll_recv_from}connect clones the connector into a boxed dial(flow.destination). The resolver is the apply’s destinations resolver and the strategy is the outbound’s address_family. Default builds one over the system resolver with default socket options; nothing in the supervisor uses it.
TCP. destination_to_socketaddrs(&dest, strategy, &resolver) resolves the destination and orders the candidates: auto keeps the resolver’s order, ipv4_only and ipv6_only filter, and prefer_ipv4 and prefer_ipv6 sort stably so the other family stays as a fallback. dialer.tcp.connect_any(&addrs) then tries the addresses in order, one at a time, each attempt bounded by DEFAULT_CONNECT_TIMEOUT (10 s). If every attempt fails, the error is ConnectionRefused failed to connect to any address (<addr>: <error>; …), listing each address with its own error. The attempts are sequential by design (not Happy Eyeballs): a dual-stack host often resolves a name to an address it cannot reach, and each failure is kept so the cause is not lost. freedom does not call set_keepalive, so a direct TCP connection keeps only the keepalive the socket policy sets (see Socket policy on every outbound socket).
UDP. dialer.udp.bind_dual(families) binds up to one socket per family. The private FreedomConnector::families maps the outbound’s address_family to the list:
address_family |
Families bound |
|---|---|
ipv4_only |
V4 |
ipv6_only |
V6 |
auto, prefer_ipv4, prefer_ipv6 |
V4 and V6 |
A family whose bind fails is left out, with a debug line udp: no V6 socket: <error> (or V4). The dial fails only when nothing binds: AddrNotAvailable udp: no usable local socket in any requested family. bind_dual called with no family at all fails with udp: no address family requested instead; families never returns an empty list, so freedom cannot meet it.
DualStackUdp is itself a DatagramLink, but that implementation sends only to IP destinations and refuses a domain with Unsupported udp: a plain dual-stack link cannot resolve a domain, because resolving a name is the application’s policy. freedom therefore wraps it in ResolvingUdp, which applies the outbound’s resolver and address_family and sends through the socket’s own poll_send_to(SocketAddr):
- A packet to an IP is sent from the socket of that IP’s family.
- A packet to a domain is sent to an address the outbound’s resolver returns for it under its
address_family. - A packet to a name that does not resolve is dropped:
poll_send_toreports it sent (Ok(buf.len())) and logsfreedom: dropping a datagram to an unresolvable <destination>at debug level, with the destination in itsDebugform. poll_recv_fromreads from whichever socket has a datagram, checking v4 before v6, and returns the source asDestination::udp(from).
A send to an address whose family has no socket fails with AddrNotAvailable udp: no local socket in the family of <peer>, for example an IPv6 address on an ipv4_only freedom. The dialers themselves are on Dialers.
Hysteria 2 and WireGuard
Section titled “Hysteria 2 and WireGuard”Both own their transport instead of dialing over a TransportConnector, and both bring it up on the first flow, not when they are built. That is why construction never fails for either, and why --test never contacts their servers.
Hysteria 2. build_outbound fills Hy2Config:
Hy2Config { server: udp(&hy2.server), // the server is a UDP endpoint server_name: hy2.server_name.clone(), password: hy2.password.expose().clone(), verify: hy2.verify.clone(), obfs: hy2.obfs.as_ref().map(obfs), max_concurrent_streams: hy2.max_concurrent_streams, socket: socket.clone(),}and builds Hy2Connector::with_address_family(config, hy2.address_family).with_resolver(dns.servers.clone()). On the first flow the connector’s slot dials the connection: it resolves the server with servers under address_family and tries each address in order. Each attempt binds its own UDP socket with UdpDialer::new(config.socket).bind(family of that address), wraps it in SalamanderSocket when obfs is set, and runs QUIC over it. A name that resolves to no usable address fails with AddrNotAvailable hysteria2: the server name resolved to no usable address, and a server none of whose addresses answers with ConnectionRefused hysteria2: no address answered (<addr>: <error>; …). Later flows of the same outbound ride the authenticated connection the slot keeps: a TCP flow as a proxy stream (OutboundStream::Hy2), a UDP flow as a UDP session over QUIC datagrams.
The CA bundle of VerifyMode::CustomCa is parsed by client_config (protocols/src/hysteria/connection.rs) each time a connection is dialed, not when the outbound is built. Its errors are all InvalidInput: hysteria2: could not read the CA file: <error>, hysteria2: the CA file holds an unusable certificate: <error> and hysteria2: the CA file contains no certificates.
A UDP flow is refused with Unsupported in two cases: hysteria2: the server does not relay UDP when the server’s authentication answer said so, and hysteria2: the server did not offer QUIC datagrams when the peer never advertised QUIC datagram support. On either, Hy2Connector::dial logs a warning, hysteria2: <error>; datagrams routed to this outbound are dropped (target etemenanki_protocols::hysteria::connector), because without it an operator routing UDP there would see nothing. The flag behind it, udp_warned, is an Arc<AtomicBool> that every clone of the connector shares, so the warning is logged once per outbound version, not once per datagram.
Hy2DatagramLink sends each datagram with its target written as an authority (format_authority). On the receive side it reassembles fragments (Defragger), skips a reply whose address does not parse or whose port is 0, and fails with BrokenPipe hysteria2: the connection behind this association is gone once the channel that delivers its replies closes. The client is covered on Hysteria 2: client.
WireGuard. build_outbound fills WgConfig:
WgConfig { private_key: *wg.private_key.expose(), peer_public_key: wg.peer_public_key, preshared_key: wg.preshared_key.as_ref().map(|k| *k.expose()), endpoint: udp(&wg.endpoint), endpoint_resolution: Some((servers.clone(), wg.endpoint_address_family)), local_addrs: wg.addresses.clone(), mtu: wg.mtu, persistent_keepalive: wg.keepalive, reserved: wg.reserved, socket: socket.clone(),}and builds WgConnector::with_address_family(config, wg.address_family).with_resolver(dns.destinations.clone()). The comment on the arm states the split: the endpoint is a proxy server like any other, so servers resolves it; what is reached inside the tunnel is a destination this process resolves, so destinations does. On the first flow acquire_device starts the device. WgDevice::start fails with InvalidInput wireguard: no tunnel-local addresses configured when local_addrs is empty. Otherwise it resolves the endpoint with servers under endpoint_address_family and takes the first address (NotFound wireguard: endpoint domain did not resolve when there is none), binds one UDP socket under the policy in that address’s family with UdpDialer::new(config.socket), connects it to the endpoint, and spawns the tunnel’s driver task. It performs no WireGuard handshake, so it succeeds even against a peer that is unreachable.
Destinations inside the tunnel are resolved with destinations by resolve_candidates("wireguard", …) (protocols/src/helpers/address_family.rs) and filtered by address_family and by the families of the tunnel-local addresses (FamilySupport::from_addrs), since the userspace netstack has no routing table to consult. A name that resolves to nothing is NotFound wireguard: destination did not resolve. When no address is left, the error is AddrNotAvailable wireguard: no usable <address_family> destination address for <host>:<port>, with (local address supports IPv4 only) or (local address supports IPv6 only) appended when the tunnel-local addresses are what ruled the candidates out; resolver errors pass through unchanged. A TCP flow tries the candidates in order, each connect inside the tunnel bounded by TCP_CONNECT_ATTEMPT_TIMEOUT (10 s); when all fail the error is TimedOut wireguard: tunnel TCP connect failed for all resolved addresses (<ip>: <error>; …). The tunnel is covered on WireGuard.
The connection and tunnel slots. Both connectors keep what they bring up in a slot shared by every clone, and bring it back after it dies:
Hysteria 2 (protocols/src/hysteria/slot.rs) |
WireGuard (protocols/src/wireguard/slot.rs) |
|
|---|---|---|
| Slot | ConnSlot behind Arc<parking_lot::Mutex<_>>, state Idle, Connecting or Ready |
DeviceSlot behind Arc<tokio::sync::Mutex<_>>, holding Option<Arc<WgDevice>> |
| Bringing it up | start_connect spawns a detached task that runs Hy2Conn::connect and writes the outcome into the slot. Dials that arrive meanwhile share its result (Shared future), so one outbound never opens two connections at once. No lock is held across the connect. |
acquire_device runs WgDevice::start when the slot holds no live device, and stores the device before it returns. |
| Detecting a dead one | Hy2Conn::is_alive is false: logged at warn as hysteria2: connection closed, reconnecting |
the driver task has ended (WgDevice::is_alive): logged at warn as wireguard: tunnel driver stopped, rebuilding |
| A failed attempt | logged at warn as hysteria2: connect failed: <error> |
logged at warn as wireguard: tunnel start failed: <error> |
| Backoff after a failure | RECONNECT_BACKOFF_BASE (1 s) doubled per consecutive failure, capped at RECONNECT_BACKOFF_MAX (30 s): 2, 4, 8, 16, then 30 s |
REBUILD_BACKOFF_BASE (1 s) and REBUILD_BACKOFF_MAX (30 s), the same schedule |
| A short-lived one | one that dies within MIN_HEALTHY_LIFETIME (10 s) of coming up counts as a failure; one that ran longer resets the count and reconnects at once |
the same, with its own MIN_HEALTHY_LIFETIME (10 s) |
| A dial during the backoff | fails with BrokenPipe hysteria2: connection is down, waiting before the next attempt |
fails with BrokenPipe wireguard: tunnel is down, waiting before the next attempt |
The short-lived rule exists because a connection that authenticates and is closed at once (a server at its user limit), or a tunnel whose start succeeds against a dead peer, would otherwise be re-dialed by every arriving flow with no backoff ever accumulating. If the Hysteria 2 connect task ends without sending its result, waiters get hysteria2: the connect task disappeared.
A reused outbound keeps what it brought up. an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply pins it: a flow opened after an apply that reuses the outbound rides the QUIC connection the first flow opened, while a changed Hysteria 2 outbound is a new version that dials a connection of its own.
The DNS outbound
Section titled “The DNS outbound”Outbound::Dns(Service) is the supervisor’s own DNS service, which intercepted port-53 flows reach when the DNS spec makes the supervisor answer. build_outbound does not build it: prepare wraps Dns::service as Target::outbound(internal_id(epoch + 1), Outbound::Dns(service)). The internal id has an empty tag, which no spec tag may have, and is rebuilt only with the DNS spec. How the service answers is on Supervisor DNS.
FanOutLink and FlowScope
Section titled “FanOutLink and FlowScope”pub const MAX_SUBS: usize = 64;
struct Sub { /// The outbound version the link is open on. id: OutboundId, /// The target the route named: the same version, or the balancer that picked it. via: OutboundId, link: MeteredDatagram<Guarded<OutboundDatagram>>,}
pub struct FanOutLink { plane: PlaneCell, ctx: FlowContext, /// The association's identity: every packet routes as this flow, toward its own target. flow: Flow, scope: FlowScope, /// Most recently sent to at the back. subs: Vec<Sub>, /// Where the last receive stopped, so every sub gets its turn. next: usize, // private open state /// The receiver, parked while no sub can deliver. recv_waker: Option<Waker>,}
impl FanOutLink { pub(crate) fn new(plane: PlaneCell, ctx: FlowContext, flow: Flow, scope: FlowScope) -> Self;}
impl DatagramLink for FanOutLink { type Addr = Destination; // poll_send_to, poll_recv_from}#[derive(Clone)]pub(crate) struct FlowScope { pub tracker: Tracker, pub session: Option<SessionId>, /// The session's wire, when the sub-links' payload counts toward it. pub charge: Option<Arc<Wire>>,}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. AppConnector::connect therefore does not dial a UDP flow at all: it returns a FanOutLink over the plane cell, the flow context, the flow and a FlowScope (the tracker, the session’s id, and the session’s wire when the connector was built with charging_datagrams). The server core sees Connected at once, and every packet is then routed on its own. See The UDP fan-out in depth.
Balancer, Member and Strategy
Section titled “Balancer, Member and Strategy”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 { /// The first healthy member in configured order. Order is the priority. Failover, /// Each healthy member in turn. RoundRobin,}
impl Strategy { pub fn parse(s: &str) -> io::Result<Self>;}
pub struct Member { pub tag: CompactString, pub target: Arc<Target>, pub probe: Destination, healthy: AtomicBool,}
impl Member { pub fn new(tag: CompactString, target: Arc<Target>, 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<Target>; pub fn select_preferring(&self, prefer: Option<&OutboundId>) -> Arc<Target>; pub fn spawn_probe( self: &Arc<Self>, tracker: &TaskTracker, token: CancellationToken, interval: Duration, timeout: Duration, resolver: Resolver, dialer: TcpDialer, );}A balancer is not an Outbound. It lives in the plane as a Target of kind TargetKind::Balancer(Arc<Balancer>) under its own tag, which shares the outbound namespace, so a route names a balancer exactly as it names an outbound. A member’s target is the same Arc<Target> the plane holds for that outbound tag in the same apply, so a member stays usable on its own under its tag, and a flow through the balancer is filed under the member’s version. The spec side is BalancerSpec { tag, members, strategy, probe_interval, probe_timeout } in supervisor/src/topology/spec_plan/route.rs. See Balancers in depth.
Data flow
Section titled “Data flow”A TCP flow through a proxy outbound
Section titled “A TCP flow through a proxy outbound”sequenceDiagram
participant RT as Server runtime
participant A as AppConnector
participant T as Target
participant O as Outbound
participant C as ProxyClient
RT->>A: connect(flow)
A->>A: admit on the session, load the plane, route
A->>T: resolve(None), a balancer picks its member now
A->>T: connect_stream(flow)
T->>O: connect_stream(flow)
O->>C: connect(flow) under the lock
C-->>O: ProxyClientConnecting
O-->>RT: boxed future, via Target and AppConnector
RT->>RT: poll: dial, codec handshake, flush
alt opened
RT->>RT: proxy_stream, then Guarded, then Metered
else any error
RT->>RT: Event ConnectFailed to the core
end
From then on the server runtime polls the stream directly. For a proxy client that means polling a ProxyClientRuntime inside the connection’s own task: no task and no channel sit between the inbound and the upstream. The routing half of this diagram is on the routing plane.
A UDP association
Section titled “A UDP association”sequenceDiagram
participant C as Server core
participant F as FanOutLink
participant P as Plane
participant T as Target
participant S as Sub-link
C->>F: Effect Open with a UDP flow, no dial
F-->>C: Event Connected at once
C->>F: poll_send_to(packet, peer)
F->>P: load, route(flow toward peer)
P-->>F: target and rule
F->>T: resolve(prefer)
alt a sub-link has this id
F->>S: poll_send_to
else no sub-link yet
F->>T: connect_datagram(flow toward peer)
T-->>F: Guarded link
F->>F: meter it as a new flow, add it to the table
F->>S: poll_send_to
end
C->>F: poll_recv_from
F->>S: poll each sub-link, starting at next
S-->>F: payload and source
F-->>C: Event Datagram from that source
The UDP fan-out in depth
Section titled “The UDP fan-out in depth”Per-packet routing
Section titled “Per-packet routing”poll_send_to(cx, buf, to) starts every packet the same way:
let flow = self.flow.toward(to.clone());let plane = self.plane.load();let routed = plane.route(&flow, &self.ctx);let rule = routed.rule;let routed = routed.target.clone();Flow::towardcopies the association’s user and source into a flow aimed at this packet’s destination and clearssniffed. 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 plane is loaded for every packet, so a route change or a new outbound version reaches the very next packet of a live association (
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow). Plane::routeintercepts port 53 for the DNS service when there is one, so a UDP DNS query becomes a sub-link on the internal DNS target.- The work per packet is
Flow::toward(a clone of the destination and of the userArc), oneArcSwapload, the first-match walk of the route table described in Routing, and linear scans of the sub-link table (at mostMAX_SUBSentries) to find the preferred and the matching sub-link, which a successful send then moves to the back.
Resolving a balancer per packet
Section titled “Resolving a balancer per packet”let prefer = self .subs .iter() .rev() .find(|s| s.via == *routed.id()) .map(|s| &s.id);let target = routed.resolve(prefer);let id = target.id();prefer is the version of the most recently used sub-link that the same route target led to. Target::resolve returns an outbound target as it is, and asks a balancer for select_preferring(prefer). The comment on this code gives the reason for resolving here rather than at open time: the sub-link is keyed by the version that carries the packet, so a member going down moves the association to another member on the next packet, instead of leaving it on a sub-link filed under the balancer.
- Under
round_robin, the preferred member is kept while it is healthy, so an association does not hop between members packet by packet (round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down). - Under
failover, the preference is ignored: the first healthy member is chosen every time, and moving back to a recovered higher member is the point (failover_ignores_the_preference). - A preference is one version of one member. After that member is rebuilt it names no member, and selection proceeds as if there were none (
a_preference_for_another_version_is_no_preference).
Sub-links keyed by OutboundId
Section titled “Sub-links keyed by OutboundId”A sub-link is found by id, the OutboundId (tag and version) of the resolved outbound. Keys are never the address of a target: outbounds outlive a single apply, one tag can have an old version beside a new one, and a freed allocation could be reused by its successor (supervisor/src/entity/id.rs). Consequences:
- Two peers routed to the same outbound version share one sub-link. For
freedom, SOCKS, Trojan, WireGuard and Hysteria 2 that is one socket or one proxied association carrying many peers. - A member reached directly by its own tag and one reached through a balancer have the same
id, so they share one sub-link. - A new version of a tag is a new key, so the next packet routed to it opens a sub-link of its own. Drains of the old version are covered on the routing plane.
blackholeis a sub-link too: aBlackholeLinkthat swallows the packets. That is how a blocked peer’s packets are dropped without harming the rest of the association.
via records the route target that led to the sub-link: the outbound itself, or the balancer that picked it. Every successful send rewrites it to the target that routed that packet, so the next packet routed there prefers this sub-link too.
Sending
Section titled “Sending”flowchart TB
A["route the packet on the current plane"] --> B["resolve a balancer, preferring the member in use"]
B --> C{"sub-link with this id?"}
C -- yes --> D["sub.poll_send_to"]
D -- "Ready Ok" --> E["move it to the back, via = routed target"]
D -- "Ready Err" --> F["remove it, drop the packet, return Ok"]
C -- no --> G["open with Target::connect_datagram"]
G -- opened --> H["meter it, add it to the table, send the packet on it"]
G -- failed --> I["log it, drop the packet, return Ok"]
The rules poll_send_to implements:
- Least recently used order.
subsis ordered by the last successful send, most recent at the back. A newly opened sub-link joins at the back before its first send. Only a successful send moves a sub-link. - A failed open drops its packet. When the open for this packet’s
idfails,poll_send_todrops the packet and returnsReady(Ok(buf.len())), so the core sees no error. The failure is logged at debug level asudp fan-out: opening <id> failed: <error>, for exampleudp fan-out: opening web@v1 failed: http carries no datagrams. - A failed send drops the sub-link and its packet. An error from a sub-link’s
poll_send_toremoves that sub-link, resetsnextto 0, logsudp fan-out: sending on <id> failed: <error>at debug level and returnsReady(Ok(buf.len())), as a lossy link would drop the packet. The next packet routed there opens a fresh sub-link. This is also how a sub-link ends when its outbound version is closed by a drain policy (Guardedfails withthe outbound this flow was opened on was drained) or when its flow is killed through the tracker. - No error reaches the core.
FanOutLink::poll_send_tonever returns an error, so the server runtime never reportsSendFailedfor a fan-out association.
Opening a sub-link and eviction
Section titled “Opening a sub-link and eviction”When no sub-link has the packet’s id, the fan-out opens one with Box::pin(target.connect_datagram(flow.clone())), and keeps with it the id, the via and the FlowMeta the sub-link will be tracked as. The FlowMeta is FlowMeta::of(&flow, session, &ctx.inbound_tag, ctx.source, id, rule, plane.epoch()): the flow toward the packet that opened the sub-link, the rule that routed that packet, and the epoch of the plane it was routed on. Because the target was already resolved, Target::connect_datagram does not consult the balancer again.
poll_opening is the only place a sub-link joins the table. When an open succeeds it:
- wraps the link with
scope.tracker.meter_datagram(meta, link, scope.charge.clone()); - pushes the new
Subto the back ofsubs; - if
subs.len() > MAX_SUBS, removes index 0, the sub-link least recently sent to (or opened), and decrementsnextwithsaturating_sub(1). Dropping the removed sub-link ends its tracked flow and drops its link; - wakes
recv_waker, so a receiver parked on an empty table polls the new sub-link.
When an open fails, it logs the failure at debug level. The table never holds more than MAX_SUBS (64) open sub-links.
Receiving
Section titled “Receiving”poll_recv_from polls every sub-link once, starting at next and wrapping around:
Ready(Ok(from)): setsnextto the following index and returnsfrom. 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 because a SOCKS control stream closed, a proxy stream failed, its version was drained withClose, or its flow was killed. It is removed,nextis reset to 0 if it fell off the end, the error is logged at debug level asudp fan-out: sub-link <id> ended: <error>, and the task wakes itself so the remaining sub-links are polled again.Pendingfrom all of them: storescx.waker()inrecv_wakerand returnsPending.
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. The fan-out therefore never ends an association because an outbound failed; the association ends when its connection ends, for example because the inbound side ended it or because the supervisor closed the session it runs under (see Users and sessions).
Each sub-link is a flow
Section titled “Each sub-link is a flow”Every sub-link is a tracked flow of its own, a MeteredDatagram<Guarded<OutboundDatagram>>:
- It is filed under the association’s session, the resolved outbound version and the rule that routed its first packet, so
sessions()and the flow table show one flow per outbound an association uses (a_udp_association_counts_each_sub_link_and_charges_its_session: three packets to two outbounds give two flows). - It counts the payload bytes sent and received on it, and applies the session’s pacing.
- When
FlowScope::chargeis set, its payload also counts toward the session’s wire bytes.serve_connection(supervisor/src/serve.rs) builds the SOCKS inbound’s connector withcharging_datagrams(), because SOCKS datagrams never cross the session’s own TCP stream. Other inbounds’ datagrams already cross what their session counts: the proxy stream, or a Hysteria 2 session’s QUIC connection. A TUN flow has no session. - Killing it through the tracker ends that sub-link alone, as any failed sub-link ends: the association runs on, and a later packet routed there opens a new flow.
Metering, pacing and kills are on Tracking; session wire bytes are on Usage.
What each sub-link does with the per-packet address
Section titled “What each sub-link does with the per-packet address”| Outbound | Each packet’s 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. The request header names the UDP command (0x03) and a 0.0.0.0:0 placeholder. |
vless, vmess |
No: the request header names one target, the destination of the packet that opened the sub-link; seal_to ignores the per-packet address, and open_from attributes every reply to that same target. |
wireguard |
Yes: each datagram goes to its own address inside the tunnel. |
hysteria2 |
Yes: each datagram carries its target, a domain as a name for the server to resolve. |
| DNS service | Answered in-process by the DNS service. |
http, shadowsocks (both families) |
No UDP: the open fails with Unsupported and the packet is dropped. |
Balancers in depth
Section titled “Balancers in depth”What may be balanced
Section titled “What may be balanced”validate (supervisor/src/build/validate.rs) checks every balancer before anything is built, so --test catches all of these. The texts are what the pinned etemenanki-app prints after configuration invalid: :
| Rule | Error |
|---|---|
| A balancer’s tag is not empty | balancer : a balancer tag must not be empty |
| No two balancers share a tag | duplicate balancer tag <tag> |
| A balancer’s tag is not also an outbound’s | duplicate balancer tag <tag> |
| A balancer has at least one member | balancer <tag> has no members |
| Every member names an outbound | balancer <tag> references unknown outbound <member> |
Every member has an upstream a TCP probe can reach (probe_target) |
balancer <tag>: outbound <member> has no upstream a TCP health probe can reach |
strategy, when the app config sets it, is exactly failover or round_robin |
balancer <tag>: unknown balancer strategy "<value>" (expected "failover" or "round_robin") |
probe_target returns upstream.server for the seven proxy clients (SOCKS, HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022) and None for the rest. freedom and blackhole have no upstream. Hysteria 2 and WireGuard listen on UDP, so a TCP connect would mark them down for ever. Members are looked up among the spec’s outbounds only, so a balancer cannot be a member of another balancer; that case gets the “unknown outbound” error.
The strategy text comes from Strategy::parse, which etemenanki-app’s lowering (app/src/lower.rs → lower_balancer) calls and wraps with balancer <tag>: . The same lowering turns an unset strategy into Strategy::Failover and unset probe_interval and probe_timeout (whole seconds) into DEFAULT_PROBE_INTERVAL and DEFAULT_PROBE_TIMEOUT. The spec itself carries Strategy and two Durations. The app’s configuration is covered on etemenanki-app: from TOML to a spec.
Probes
Section titled “Probes”commit starts probing for every balancer the apply built, once the new plane is published:
let token = self.root.child_token();balancer.spawn_probe( &self.shared.tracker, token.clone(), interval, timeout, dns.servers.clone(), TcpDialer::new(self.socket.clone()),);self.probes.insert(tag, token);The tasks run on the supervisor’s TaskTracker, under a child of the supervisor’s root token that Actor::probes keeps by balancer tag. They resolve with servers, the resolver the members’ upstreams are dialed with, and connect under the socket policy.
spawn_probe spawns one task per member:
stateDiagram-v2 [*] --> Probing Probing --> Recording: connected, or failed, or timed out Probing --> [*]: token cancelled, probe abandoned Recording --> Sleeping: set_healthy, log if the state changed Sleeping --> Probing: interval elapsed Sleeping --> [*]: token cancelled
probe(dest, timeout, resolver, dialer)resolvesmember.probewithAddressFamilyStrategy::Autoand triesdialer.connecton each address in turn. It is up if any address accepts; the connection is dropped at once.tokio::time::timeoutbounds the whole probe, resolution included, bytimeout; each connect attempt is also bounded by the dialer’sDEFAULT_CONNECT_TIMEOUT. A resolution error, no address, or the timeout count as down.- The probe is a TCP connect and nothing more. It uses neither the member’s transport, nor its credentials, nor its
address_family: it answers “is the upstream reachable”, which is the failure a balancer exists to route around. Anything richer would need credentials and would make the probe a second implementation of the proxy client. set_healthyswaps the member’sAtomicBoolwithOrdering::Relaxedand returns the old value. A change is logged at info level asbalancer member <tag> is now uporbalancer member <tag> is now down.- Members start healthy.
Member::newsetshealthytotrue, and the first probe runs as soon as the task starts. An upstream that has never been probed counts as healthy, because starting every member down would drop all traffic for a whole interval after every start. - Cancellation abandons a probe in flight. Both the probe and the sleep race
token.cancelled()in atokio::select!, so a cancel ends the task at once instead of waiting a probe out.
The end-to-end test traffic_moves_off_a_member_that_stops_answering runs two app instances with probe_interval = 1. The first member accepts TCP but is not a proxy, so it probes healthy while flows through it fail; once it stops accepting, the probe marks it down and flows move to the working member.
Selection
Section titled “Selection”pub fn select(&self) -> Arc<Target> { let mut healthy = self.members.iter().filter(|m| m.is_healthy()); let chosen = match self.strategy { Strategy::Failover => healthy.next(), Strategy::RoundRobin => match healthy.clone().count() { 0 => None, n => healthy.nth(self.next.fetch_add(1, Ordering::Relaxed) % n), }, }; chosen.unwrap_or(&self.members[0]).target.clone()}select_preferring(prefer) returns the preferred member’s target when the strategy is RoundRobin, prefer is Some, and a member whose target.id() equals it is healthy. Otherwise it calls select.
Who calls them:
- A TCP flow:
AppConnector::connectcallsTarget::resolve(None)once, and the flow is filed and guarded under the chosen member.Target::connect_streamresolves again, which for an outbound target is the target itself. - A UDP packet:
FanOutLink::poll_send_tocallsTarget::resolve(prefer)for every packet, as described above.
Selection never runs at route time, so a member going down affects the next flow or packet, not the next apply.
flowchart TB
P{"round_robin, and the preferred member is healthy?"} -- yes --> K["keep the preferred member"]
P -- no --> A["members whose healthy flag is set"]
A --> B{"any healthy?"}
B -- no --> F["first configured member"]
B -- yes --> C{"strategy"}
C -- failover --> D["first healthy member in configured order"]
C -- round_robin --> E["healthy member at next.fetch_add(1) mod 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; a UDP association keeps its member while it stays healthy. |
| 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: a flow sent at an upstream that may have recovered is better than one dropped for certain. The alternative is a balancer that drops every flow from the moment each member’s latest probe has failed until one of them succeeds again, up to a whole probe_interval after an upstream has recovered. Selection is lock-free: relaxed atomic loads of the members’ health flags, plus one relaxed fetch_add under round_robin when any member is healthy. failover stops loading at the first healthy member; round_robin makes up to two passes, one to count the healthy members (healthy.clone().count()) and one to reach the chosen one (healthy.nth(k)).
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, which is why health is established out of band, as in Xray’s observatory.
Balanced flows and drains
Section titled “Balanced flows and drains”A flow through a balancer is opened on the member’s Target and guarded by the member version’s close token, never the balancer’s. a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers pins it: closing the balancer’s own target leaves the stream open, and closing the member’s target aborts it. The plan emits Drain steps for outbound versions only, so a balancer is never drained itself.
State across applies
Section titled “State across applies”A balancer holds its members’ targets, so plan rebuilds it with any of them:
let members_kept = !balancer.members.iter().any(|m| rebuilt.contains(m));match running_version { Some(version) if prev == Some(balancer) && members_kept => { // Step::Reuse(Resource::Balancer(tag)) at the running version } _ => { // Step::Build(Resource::Balancer(tag)) at the next version }}| Change | What happens to the balancer |
|---|---|
| Neither the balancer’s spec nor any member changed | Step::Reuse. The running Arc<Target> is kept, and with it the Balancer: learned health, the round-robin counter and the probe tasks, whose token commit keeps (probes.retain keeps every tag the plan reuses). |
| The balancer’s spec changed, or any member was rebuilt | Step::Build at the next version. The new Balancer starts with every member healthy and gets new probe tasks under a new child token; the old token is cancelled in the same commit. |
| The DNS spec changed | Every member is rebuilt (every probeable outbound uses DNS), so every balancer is rebuilt and relearns its health. |
| The balancer was removed | The plane is republished without it and its probe token is cancelled. Its version is kept, so a balancer added back under the same tag gets version 2 (removing_a_balancer_alone_republishes_the_plane). |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is pins the first row and the member half of the second; a change to the balancer’s own spec has no dedicated test. Starting and stopping probes happens inside the plane publication in commit, which every balancer build or removal triggers. At shutdown, Actor::shutdown cancels every probe token before it closes the task tracker and waits on it.
Socket policy on every outbound socket
Section titled “Socket policy on every outbound socket”The socket policy is a SocketOptions set once on the builder with SupervisorBuilder::socket_options. It holds for the supervisor’s lifetime; an apply cannot change it. Its doc comment lists what it covers: each dial of an outbound or a balancer probe, each UDP socket of a direct flow, a SOCKS UDP relay, a Hysteria 2 or WireGuard tunnel, and every query a resolver sends from this host. This is where a mobile VPN keeps the proxy’s own traffic out of its tunnel.
| Socket | Opened by | How the policy reaches it |
|---|---|---|
freedom TCP connect |
TcpDialer::connect_any |
Dialer::new(socket) in build_outbound |
freedom UDP sockets, one per family |
UdpDialer::bind_dual |
the same Dialer |
Proxy upstream TCP: HTTP, SOCKS (CONNECT and the association’s control stream), Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022 |
TransportConnector::dial |
Dialer::new(socket) in transport |
| SOCKS UDP relay socket | the bind closure of SocksUdpLink::associate |
UdpDialer::new(socket) in build_outbound |
| Hysteria 2 QUIC socket | bind_socket in protocols/src/hysteria/connection.rs |
Hy2Config.socket |
| WireGuard tunnel socket | WgDevice::start |
WgConfig.socket |
| Balancer probe connects | probe |
TcpDialer::new(self.socket) in commit |
| Queries to configured DNS servers | the resolver’s upstreams | ResolverOptions.socket in build_dns (Supervisor DNS) |
| DNS queries sent through the proxy | RoutedDialer, over the route table |
the socket of whichever outbound carries them |
The system resolver (getaddrinfo) |
the C library | not reached: it opens its sockets where no policy can |
SocketOptions::apply runs between creating a socket and binding or connecting it, the only window in which most options can be set. It sets the interface, then the mark (SO_MARK), then runs the hook. TcpDialer::socket then sets the policy’s tcp_keepalive idle time when it has one, and binds bind_address if it is of the socket’s family. On proxy upstreams that keepalive does not last: after the connect, TransportConnector::dial calls set_keepalive, which replaces it with the fixed 120 s, 30 s and 3-probe schedule (see The transport under a proxy client). A freedom TCP connection keeps only what the policy sets. UdpDialer::bind makes an IPv6 socket v6-only before apply, sets the socket non-blocking after it, then binds bind_address of the same family or the wildcard address, on an ephemeral port. A platform without a requested option fails with Unsupported instead of silently ignoring it. The full policy is on Dialers.
A hook that returns an error fails that socket, and with it the dial: nothing leaves on a socket the policy did not accept. For freedom TCP the hook’s error lands in failed to connect to any address (…), and the SOCKS client sees a failure reply (a_refusing_hook_fails_the_dial). A UDP family whose socket the hook refuses is left out of bind_dual, and the dial fails only when no family remains. The mobile bindings (ffi/src/platform.rs → socket_options) install a hook that asks the platform to protect each socket’s file descriptor, and fails with PermissionDenied the platform refused to protect the socket when it declines.
check builds the supervisor it validates with SocketOptions::default(), and construction opens no socket, so --test never runs a hook.
Adding an outbound protocol
Section titled “Adding an outbound protocol”- Add a variant to
OutboundProtocolSpec(supervisor/src/topology/spec_plan/outbound.rs). A proxy client over a stream transport takes aProxyUpstream { server, transport }. The spec is compared with==by the planner, so every field that changes behaviour must be in it. - Add its rules to
validate_outbound(supervisor/src/build/validate.rs), includingoutbound_transportwhen it has aProxyUpstream. Decide whatprobe_targetreturns: a TCP upstream makes it balanceable,Nonerefuses it as a member. - Check
plan’suses_dnsmatch: every outbound butblackholeis rebuilt when the DNS spec changes. A protocol that resolves nothing can joinblackholethere. - Add an
Outboundvariant (supervisor/src/topology/outbound/mod.rs) and both arms. A stream goes throughProxyClientandproxy_stream(it becomesOutboundStream::Proxy), or gets a variant of its own inOutboundStreamif it should not be boxed. A datagram link is boxed withOutboundDatagram::new; a protocol without UDP returnsno_udp("<name>"). - For a
ProxyClient, add aBUFconstant that holds the largest frame the protocol opens and exceeds the codec’sSTAGING_RESERVE, orProxyClientRuntime::newpanics on the first flow. - Add the
build_outboundarm. Pass the socket policy to every socket it can open (Dialer::new(socket),UdpDialer::new(socket), or asocketfield in the protocol’s config), resolve upstream names withdns.serversand destinations it reaches itself withdns.destinations, and derive keys once, outside themakeclosure. - Teach the front ends to lower it: etemenanki-app’s
app/src/lower.rs, and katana’s lowering if a panel can express it (katana inbounds and outbounds). If subscribe files may name it, add an arm tonode_outboundinapp/src/subscribe.rs, which turns each subscribe node into theOutboundConfigan[[outbound]]table would give, so the node then goes through the same lowering; the subscribe format is described on Subscribe files. - Test it: a
socket_policy.rscase for each new kind of socket, ahot_swap.rscase if it holds a connection across flows, andvalidate.rscases for every new rule.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| A UDP association is routed per packet, on the plane current for that packet | AppConnector::connect returns a FanOutLink for UDP flows; poll_send_to loads the plane and routes every packet |
one_association_routes_each_peer_separately (e2e_udp_route.rs), a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow (hot_swap.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 its source; the fan-out passes it through |
replies_from_several_peers_merge_back_correctly (e2e_udp_route.rs) |
| UDP rules see the UDP network | The per-packet flow keeps DialNetwork::Udp |
network_separates_tcp_from_udp (e2e_route_context.rs) |
| Each outbound an association uses is one tracked flow | poll_opening meters each sub-link with meter_datagram |
a_udp_association_counts_each_sub_link_and_charges_its_session (tracking.rs) |
At most MAX_SUBS sub-links are open per association |
poll_opening removes index 0 after a push past MAX_SUBS |
not tested directly |
| A failed open or send drops its packet and never fails the association | poll_send_to returns Ready(Ok(buf.len())) on both |
not tested directly |
| Sub-link keys never alias | Keys are OutboundIds, tag and version, never addresses |
by construction |
The ProxyClient lock is never held across I/O |
ProxyClient::connect returns the dial future; it is polled outside the lock |
by construction |
| A proxy dial resolves only after the upstream dial, the codec handshake and the flush have completed | ProxyClientConnecting polls the runtime until poll_connected is ready |
server_runtime_relays_through_a_client_runtime_in_one_task, upstream_dial_failure_is_connect_failed_not_connected, refused_upstream_handshake_is_connect_failed_not_connected (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 refuse it | connect_datagram returns Unsupported for http, shadowsocks and shadowsocks-2022 |
not tested directly |
| Outbound sockets are opened under the socket policy | Every dialer in build_outbound and every probe dialer is built from self.socket |
direct_flows_are_dialed_on_hooked_sockets, a_socks_upstream_is_dialed_on_hooked_sockets (socket_policy.rs); Hysteria 2, WireGuard and probe sockets are not covered |
| A refusing hook fails the dial | SocketOptions::apply returns the hook’s error before any connect |
a_refusing_hook_fails_the_dial (socket_policy.rs) |
| An unchanged outbound keeps what it holds across an apply | prepare reuses the running Arc<Target> |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply (hot_swap.rs) |
A changed outbound is built as a new version that new flows use, and a Close drain closes the flows on the old one |
plan builds the next version and emits a Drain step for the old one; Guarded fails once its version is closed |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close (hot_swap.rs, a freedom outbound) |
| A balanced flow drains with its member, not with the balancer | Target::resolve picks the member before the flow is guarded |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers (plane.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 (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[0] |
every_member_down_still_selects_rather_than_dropping |
round_robin keeps a healthy preferred member and drops it once down |
select_preferring |
round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down |
failover ignores a preference |
select_preferring checks the strategy first |
failover_ignores_the_preference |
| A preference for another version is no preference | select_preferring compares OutboundIds |
a_preference_for_another_version_is_no_preference |
| A balancer is never empty | validate (EmptyBalancer) and Balancer::new |
a_balancer_needs_a_member (validate.rs), a_balancer_needs_at_least_one_member (balancer.rs) |
| Strategy names are exact | Strategy::parse accepts only failover and round_robin |
strategy_names_are_validated |
| Only outbounds with a TCP upstream can be balanced | probe_target returns None for freedom, blackhole, Hysteria 2 and WireGuard |
a_balancer_member_needs_an_upstream_a_tcp_probe_can_reach (validate.rs), a_member_with_no_upstream_is_refused (e2e_balancer.rs) |
| Members are outbounds, never balancers | Members are looked up among the spec’s outbounds only | a_balancer_member_may_not_be_another_balancer (validate.rs) |
| A balancer tag is not an outbound tag | validate refuses the overlap as a duplicate |
a_balancer_may_not_share_an_outbounds_tag (validate.rs) |
| A balancer keeps its health unless it or a member is rebuilt | plan reuses it when members_kept |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is (plan.rs) |
| A member that stops answering stops receiving flows | The probe marks it down; select skips it |
traffic_moves_off_a_member_that_stops_answering (e2e_balancer.rs) |
| Probes stop with their balancer | Each probe task selects on its balancer’s token, cancelled on rebuild, removal and shutdown | not tested directly |
Failure paths and cancellation
Section titled “Failure paths and cancellation”Errors
Section titled “Errors”| Where | Condition | Kind | Message |
|---|---|---|---|
prepare |
Any construction error below | none (wrapper) | building outbound <tag>@v<n> failed: <error> |
transport |
TLS shape without TLS material, WebSocket without host, gRPC without authority (refused by validation first) | InvalidInput |
missing tls, missing ws host, missing grpc authority |
transport → ClientConfig::with_verify_mode |
A CA bundle that does not parse | Other |
OpenSSL’s error text |
transport → ClientConfig::with_verify_mode |
A CA bundle with no certificate | InvalidInput |
no certificate in CA PEM bundle |
TransportConnector::dial |
A UDP destination (never passed by the supervisor) | Unsupported |
a proxy transport carries no datagrams of its own |
Balancer::new |
No members (refused by validation first) | InvalidInput |
a balancer needs at least one outbound, wrapped as building balancer <tag> failed: … |
Strategy::parse |
Unknown strategy name | InvalidInput |
unknown balancer strategy "<value>" (expected "failover" or "round_robin") |
Target::connect_stream, connect_datagram |
A member that is itself a balancer (made unreachable by validation) | Other |
balancer member <id> is itself a balancer |
connect_datagram |
http, shadowsocks, shadowsocks-2022 |
Unsupported |
http carries no datagrams, shadowsocks carries no datagrams, shadowsocks-2022 carries no datagrams |
connect_stream, connect_datagram |
freedom, wireguard or hysteria2 produced the other kind |
Other |
freedom dialed a TCP flow as UDP, freedom dialed a UDP flow as TCP, and the same with wireguard and hysteria2 |
proxy_stream, proxy_datagram |
A client produced the other kind | Other |
a TCP flow was dialed as datagrams, a UDP flow was dialed as a stream |
ProxyClientRuntime |
The upstream closed during the handshake | UnexpectedEof |
upstream closed during the handshake |
ProxyClientRuntime |
The upstream closed with part of a frame unread | UnexpectedEof |
upstream closed inside a frame |
ProxyClientRuntime |
An upstream frame larger than BUF |
InvalidData |
upstream frame larger than the client runtime's buffer |
ProxyClientRuntime |
A datagram whose sealed frame does not fit the staging buffer | InvalidInput |
frame larger than the client runtime's buffer |
ProxyClientRuntime |
The codec rejected the upstream’s bytes | InvalidData |
the codec’s own error |
ProxyClientRuntime |
A codec broke its contract | Other, InvalidInput |
codec opened an impossible frame, codec consumed more reply than it was given, codec handshake step made no progress, codec sealed an impossible byte count, codec refused a packet that fits its reserve |
ProxyClientRuntime |
The transport yielded a datagram socket | Unsupported |
proxy client runtime needs a stream to the upstream, got a datagram socket |
ProxyClientRuntime |
Any call after the wire failed | NotConnected |
none |
freedom TCP, resolution |
The name resolved to nothing | NotFound |
dial: destination did not resolve |
freedom TCP, resolution |
No address of an allowed family | AddrNotAvailable |
dial: no usable <strategy> destination address for <host>:<port> |
freedom TCP, connect |
One address does not answer in time | TimedOut |
connect to <addr> timed out (collected, not returned alone) |
freedom TCP, connect |
No address connects | ConnectionRefused |
failed to connect to any address (<addr>: <error>; …) |
freedom UDP |
No family binds | AddrNotAvailable |
udp: no usable local socket in any requested family |
UdpDialer::bind_dual |
No family requested (never from freedom) |
AddrNotAvailable |
udp: no address family requested |
DualStackUdp as a DatagramLink |
A domain destination (freedom wraps it in ResolvingUdp) |
Unsupported |
udp: a plain dual-stack link cannot resolve a domain |
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 | The packet is dropped with a debug line |
SocksUdpLink::associate |
Method refused, account refused | PermissionDenied |
auth method not supported, server rejects account |
SocksUdpLink::associate |
Non-zero reply status | ConnectionRefused |
server rejects request: <status> |
SocksUdpLink::associate |
The relay is named by domain | Unsupported |
socks: the relay address is a domain |
SocksUdpLink::associate |
The server closed mid-handshake | UnexpectedEof |
socks: server closed during the handshake |
SocksUdpLink::associate |
A reply that does not parse | as the parser returns it | for example unexpected server version |
SocksUdpLink |
The control stream reached EOF, or a call after a control read error | BrokenPipe |
socks: the control connection closed |
| Hysteria 2 dial | The server name leaves no usable address | AddrNotAvailable |
hysteria2: the server name resolved to no usable address |
| Hysteria 2 dial | No address completes a connection | ConnectionRefused |
hysteria2: no address answered (<addr>: <error>; …) |
| Hysteria 2 dial | The CA bundle is unreadable, holds a bad certificate, or holds none | InvalidInput |
hysteria2: could not read the CA file: <error>, hysteria2: the CA file holds an unusable certificate: <error>, hysteria2: the CA file contains no certificates |
| Hysteria 2 dial | The slot is backing off | BrokenPipe |
hysteria2: connection is down, waiting before the next attempt |
| Hysteria 2 dial | The connect task ended without a result | Other |
hysteria2: the connect task disappeared |
| Hysteria 2 UDP open | The server does not relay UDP, or offered no QUIC datagrams | Unsupported |
hysteria2: the server does not relay UDP, hysteria2: the server did not offer QUIC datagrams |
Hy2DatagramLink receive |
The channel delivering its replies closed | BrokenPipe |
hysteria2: the connection behind this association is gone |
WgDevice::start |
No tunnel-local address | InvalidInput |
wireguard: no tunnel-local addresses configured |
WgDevice::start |
The endpoint resolves to nothing | NotFound |
wireguard: endpoint domain did not resolve |
| WireGuard dial | The slot is backing off | BrokenPipe |
wireguard: tunnel is down, waiting before the next attempt |
| WireGuard dial, resolution | The destination resolves to nothing | NotFound |
wireguard: destination did not resolve |
| WireGuard dial, resolution | No address of an allowed family | AddrNotAvailable |
wireguard: no usable <strategy> destination address for <host>:<port>, with (local address supports IPv4 only) or (local address supports IPv6 only) when the tunnel-local addresses ruled them out |
| WireGuard TCP dial | No address connects inside the tunnel | TimedOut |
wireguard: tunnel TCP connect failed for all resolved addresses (<ip>: <error>; …) |
Guarded |
The flow’s outbound version was closed by a drain policy | ConnectionAborted |
the outbound this flow was opened on was drained |
FanOutLink |
An open, a send or a receive of a sub-link fails | none | Logged; the packet is dropped or the sub-link removed |
Log lines
Section titled “Log lines”| Level | Target | Text |
|---|---|---|
| debug | etemenanki_supervisor::topology::outbound::udp_fanout |
udp fan-out: opening <id> failed: <error> |
| debug | etemenanki_supervisor::topology::outbound::udp_fanout |
udp fan-out: sending on <id> failed: <error> |
| debug | etemenanki_supervisor::topology::outbound::udp_fanout |
udp fan-out: sub-link <id> ended: <error> |
| debug | etemenanki_supervisor::topology::outbound::freedom |
freedom: dropping a datagram to an unresolvable <destination> |
| debug | etemenanki_environment::dial::udp |
udp: no V4 socket: <error>, udp: no V6 socket: <error> |
| debug | etemenanki_protocols::transports::keepalive |
could not enable TCP keepalive: <error> |
| info | etemenanki_supervisor::topology::balancer |
balancer member <tag> is now up, balancer member <tag> is now down |
| warn | etemenanki_protocols::hysteria::connector |
hysteria2: <error>; datagrams routed to this outbound are dropped (once per outbound version) |
| warn | etemenanki_protocols::hysteria::slot |
hysteria2: connect failed: <error>, hysteria2: connection closed, reconnecting |
| warn | etemenanki_protocols::wireguard::slot |
wireguard: tunnel start failed: <error>, wireguard: tunnel driver stopped, rebuilding |
<id> is the sub-link’s OutboundId. The Hysteria 2 and WireGuard drivers log more lines of their own; they are on Hysteria 2: client and WireGuard.
How callers see each outcome
Section titled “How callers see each outcome”- An
Errfrom a TCP flow’s dial future becomesEvent::ConnectFailedin the server core. The SOCKS inbound answers it with a failure reply. - An
Errfrom aGuardedorMeteredstream after it opened is an outbound error for that key in the server runtime (The server runtime). - The fan-out swallows open failures, send failures and receive errors, so a UDP association never sees an outbound error from it. The association ends with its connection, as described under Receiving.
- A construction error refuses the apply and changes nothing.
Cancellation
Section titled “Cancellation”Cancellation in this layer is by drop, except for the probes and the Hysteria 2 connect task:
- Dropping the
StreamFutureorDatagramFutureof a proxy client before it resolves drops the transport connection and the partial handshake; for a SOCKS association it drops the control stream too. Forfreedomit abandons the lookup and the connect attempt. - Dropping a Hysteria 2 dial does not cancel its connection attempt:
start_connectrunsHy2Conn::connectin a task of its own, which completes even if every waiter is dropped and leaves the connection in the outbound’s slot for the next flow (Hysteria 2: client). - Dropping a WireGuard dial after
WgDevice::starthas returned leaves the started tunnel in the outbound’s slot, driver task included (WireGuard). - Dropping an
OutboundStreamcloses the TCP stream, or drops the client runtime and with it the upstream stream. - Dropping a
FanOutLinkdrops every sub-link and any sub-link still being opened. It happens when the connection’s runtime ends, and each sub-link’s tracked flow ends as it drops. - Probe tasks are cancelled through their balancer’s token when the balancer is rebuilt or removed, and at shutdown. A probe in flight is abandoned, not waited out.
- The supervisor’s outbound code spawns only the probe tasks. The protocols crate spawns the Hysteria 2 connect task and the drivers of the Hysteria 2 connection and the WireGuard tunnel, described on their protocol pages.
Limits
Section titled “Limits”| Constant | Where | Value | Meaning |
|---|---|---|---|
MAX_SUBS |
supervisor/src/topology/outbound/udp_fanout.rs |
64 | Open sub-links per UDP association; the least recently sent to is closed first |
DEFAULT_PROBE_INTERVAL |
supervisor/src/topology/balancer.rs |
30 s | Sleep between two probes of one member, used by the app when probe_interval is unset |
DEFAULT_PROBE_TIMEOUT |
supervisor/src/topology/balancer.rs |
5 s | Bound on one probe, resolution and every connect attempt included, used by the app when probe_timeout is unset |
DEFAULT_CONNECT_TIMEOUT |
environment/src/dial/tcp.rs |
10 s | One TCP connect attempt by TcpDialer (see Dialers) |
TCP_KEEPALIVE_IDLE, TCP_KEEPALIVE_INTERVAL, TCP_KEEPALIVE_RETRIES |
protocols/src/transports/keepalive.rs |
120 s, 30 s, 3 | Keepalive schedule TransportConnector::dial sets on every proxy upstream TCP connection |
RECONNECT_BACKOFF_BASE, RECONNECT_BACKOFF_MAX |
protocols/src/hysteria/slot.rs |
1 s, 30 s | Backoff between Hysteria 2 connection attempts after consecutive failures |
REBUILD_BACKOFF_BASE, REBUILD_BACKOFF_MAX |
protocols/src/wireguard/slot.rs |
1 s, 30 s | Backoff between WireGuard tunnel starts after consecutive failures |
MIN_HEALTHY_LIFETIME |
protocols/src/hysteria/slot.rs, protocols/src/wireguard/slot.rs |
10 s | A connection or tunnel lost sooner counts as a failed attempt |
TCP_CONNECT_ATTEMPT_TIMEOUT |
protocols/src/wireguard/slot.rs |
10 s | One TCP connect inside a WireGuard tunnel, per resolved address |
RECV_BUF |
protocols/src/socks/udp_link.rs |
64 KiB | Largest SOCKS relay packet read; one boxed buffer per association |
sink |
protocols/src/socks/udp_link.rs |
256 bytes | Buffer the SOCKS control stream is drained into |
scratch |
protocols/src/socks/udp_link.rs |
2,048 bytes | Initial capacity of the reused buffer a SOCKS UDP packet is built in |
HTTP_BUF, SOCKS_BUF, TROJAN_BUF, VLESS_BUF |
supervisor/src/topology/outbound/mod.rs |
16,384 bytes | Client runtime buffer size; two per flow |
VMESS_BUF |
supervisor/src/topology/outbound/mod.rs |
32,768 bytes | VMess client runtime buffer size |
SS_BUF |
supervisor/src/topology/outbound/mod.rs |
20,480 bytes | SsStream::BUF_SIZE |
SS2022_BUF |
supervisor/src/topology/outbound/mod.rs |
65,569 bytes | Ss2022Stream::BUF_SIZE = MAX_RECORD_LEN |
Per UDP association, MAX_SUBS bounds the number of open sub-links. A Trojan, VLESS or VMess sub-link is an upstream connection of its own through the transport, with two client buffers; a SOCKS sub-link is a control connection of its own plus a relay socket and its 64 KiB receive buffer. The workspace-wide limits table is on Limits, timeouts and memory.
The supervisor’s unit tests are #[path] modules of the library under supervisor/tests/unit/; its integration tests (socket_policy, hot_swap, tracking) run supervisors in-process against loopback echo servers. The app’s end-to-end tests are one integration crate, app/tests/integration.rs, that spawns the real binary. The SocksUdpLink tests are in the protocols pipeline crate.
cargo test -p etemenanki-supervisor --lib topology::balancercargo test -p etemenanki-supervisor --lib topology::planecargo test -p etemenanki-supervisor --lib topology::spec_plan::plancargo test -p etemenanki-supervisor --lib build::validatecargo test -p etemenanki-supervisor --test socket_policycargo test -p etemenanki-supervisor --test hot_swapcargo test -p etemenanki-supervisor --test tracking a_udp_associationcargo test -p etemenanki-app --test integration e2e_udp_routecargo test -p etemenanki-app --test integration e2e_balancercargo test -p etemenanki-protocols --test pipeline pipeline::socks::| Test | File | Behaviour pinned |
|---|---|---|
failover_takes_the_first_healthy_member_in_order |
supervisor/tests/unit/balancer.rs |
With a down, b is picked; once a recovers, a is picked again. |
round_robin_cycles_only_through_healthy_members |
supervisor/tests/unit/balancer.rs |
With b down, four picks give a, c, a, c. |
every_member_down_still_selects_rather_than_dropping |
supervisor/tests/unit/balancer.rs |
All down: the first member is still returned. |
round_robin_keeps_a_healthy_preferred_member_and_drops_it_once_down |
supervisor/tests/unit/balancer.rs |
Preferring b gives b four times; with b down, every pick is a or c. |
failover_ignores_the_preference |
supervisor/tests/unit/balancer.rs |
Preferring b under failover still gives a. |
a_preference_for_another_version_is_no_preference |
supervisor/tests/unit/balancer.rs |
A preference for version 7 of b alternates a, b, a, b. |
a_balancer_needs_at_least_one_member |
supervisor/tests/unit/balancer.rs |
Balancer::new refuses an empty list. |
strategy_names_are_validated |
supervisor/tests/unit/balancer.rs |
failover and round_robin parse; roundrobin does not. |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers |
supervisor/tests/unit/plane.rs |
Closing the balancer’s target leaves a balanced stream open; closing the member’s aborts it. |
closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive |
supervisor/tests/unit/plane.rs |
A drained datagram link fails its send and a receive already waiting, which is how a drained fan-out sub-link ends. |
a_balancer_is_rebuilt_exactly_when_one_of_its_members_is |
supervisor/tests/unit/plan.rs |
Changing a non-member keeps the balancer at v1; changing a member builds v2 and republishes the plane. |
removing_a_balancer_alone_republishes_the_plane |
supervisor/tests/unit/plan.rs |
Removing only a balancer publishes a plane, drains nothing, and keeps its version. |
a_balancer_needs_a_member, a_balancer_member_needs_an_upstream_a_tcp_probe_can_reach, a_balancer_member_must_be_a_known_outbound, a_balancer_member_may_not_be_another_balancer, a_balancer_may_not_share_an_outbounds_tag |
supervisor/tests/unit/validate.rs |
The balancer rules in What may be balanced; the probe rule is checked for freedom, blackhole, Hysteria 2 and WireGuard. |
direct_flows_are_dialed_on_hooked_sockets |
supervisor/tests/socket_policy.rs |
A direct TCP dial is one hooked stream socket; a direct UDP flow uses at least one hooked datagram socket. |
a_socks_upstream_is_dialed_on_hooked_sockets |
supervisor/tests/socket_policy.rs |
Through a SOCKS upstream: the CONNECT dial and the association’s control connection are two hooked stream sockets, and the relay socket is one hooked datagram socket. |
a_refusing_hook_fails_the_dial |
supervisor/tests/socket_policy.rs |
A hook that returns PermissionDenied makes the SOCKS request fail. |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
supervisor/tests/hot_swap.rs |
After the default route becomes block, the next packet of a live association is dropped, while an open TCP flow keeps echoing. |
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close |
supervisor/tests/hot_swap.rs |
Version 1 of the freedom outbound direct is drained with Keep and its flow runs on; version 2 is drained with Close and its flow is closed; new flows use the newest version. |
an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply |
supervisor/tests/hot_swap.rs |
A reused Hysteria 2 outbound adds no server session; a changed one dials a second, and the first flow keeps working. |
a_udp_association_counts_each_sub_link_and_charges_its_session |
supervisor/tests/tracking.rs |
Three packets to two outbounds are two flows with exact byte counts, and the payload is charged to the SOCKS session’s wire. |
one_association_routes_each_peer_separately |
app/tests/integration/e2e_udp_route.rs |
One SOCKS association: 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, flows leave a first member once it stops accepting TCP. |
a_member_with_no_upstream_is_refused |
app/tests/integration/e2e_balancer.rs |
--test rejects a balancer over a freedom outbound. |
server_runtime_relays_through_a_client_runtime_in_one_task, upstream_dial_failure_is_connect_failed_not_connected, refused_upstream_handshake_is_connect_failed_not_connected |
concepts/tests/client.rs |
A client runtime is a server runtime’s outbound in one task, and a failed dial or refused handshake is ConnectFailed. |
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 |
A well-formed reply from a socket other than the relay is dropped; the relay’s own reply that follows is returned. |
udp_link_on_a_dual_stack_socket_hears_an_ipv4_relay |
protocols/tests/pipeline/socks.rs |
Linux only: a link bound on [::]:0 hears an IPv4 relay’s replies from the IPv4-mapped address. |
app_client_ws_xray_server_plain and the other app_client_* tests |
app/tests/integration/e2e_xray.rs |
The VLESS client against a real Xray server over WebSocket, gRPC and TLS; skipped without go. |
app_client_vmess_ws_xray_server_early_data_plain |
app/tests/integration/e2e_xray_vmess.rs |
The VMess client against Xray over WebSocket with early data; skipped without go. |
app_socks_to_hysteria2_outbound, app_socks_to_hysteria2_with_salamander, app_socks_udp_to_hysteria2_outbound |
app/tests/integration/e2e_hysteria.rs |
TCP, Salamander and UDP through the Hysteria 2 outbound against the upstream server built from the vendored tree; skipped without go. |
app_socks_to_wireguard_outbound_tcp |
app/tests/integration/e2e_wg.rs |
TCP through the WireGuard outbound to an in-test peer. |
Several fan-out properties have no dedicated test: the MAX_SUBS eviction, the drop on a failed open or send, and the per-packet balancer preference inside FanOutLink (the preference itself is unit-tested on Balancer). The socket policy is not tested on Hysteria 2, WireGuard or probe sockets. If you change FanOutLink, a unit test over a hand-built Plane with blackhole and scripted targets is the missing piece; see Testing.