Connector and UDP fan-out
Source files: 24 · checked against katana v3.0.1 · Etemenanki 596916d
katana/src/connector.rskatana/src/router.rskatana/src/rule.rskatana/src/meter.rskatana/src/serve.rskatana/src/manager/transport.rskatana/src/outbound/mod.rskatana/src/outbound/proxy.rskatana/src/outbound/freedom.rskatana/tests/unit/connector.rskatana/tests/unit/meter.rskatana/tests/unit/rule.rskatana/tests/unit/serve.rskatana/tests/unit/e2e.rskatana/tests/integration/sniff.rsEtemenanki/concepts/src/link.rsEtemenanki/concepts/src/core.rsEtemenanki/concepts/src/runtime.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/helpers/address_family.rsEtemenanki/protocols/src/vless/codec.rsEtemenanki/protocols/src/vmess/codec.rsEtemenanki/environment/src/routing.rsEtemenanki/environment/src/dial/udp.rs
Every flow an inbound decodes on a katana node ends up in one call: KatanaConnector::connect. A TCP request, each mux sub-flow, a UDP association and every Hysteria 2 stream arrive there as a Flow<UserTag>. The connector decides, synchronously and in a fixed order, whether the user may open it, which outbound carries it, whether an audit rule forbids it, and how its bytes are billed. It hands the per-connection runtime either a metered byte stream or a FanOut, the datagram link that routes a UDP association one packet at a time.
This page is for contributors who change routing, audit, billing or outbound code. It walks through src/connector.rs step by step, then covers FanOut in depth, and ends with ResolvingUdp, the direct outbound’s UDP socket. Admission itself, the token bucket and the server runtime that drives the connector each have their own page.
Responsibilities
Section titled “Responsibilities”| Concern | Owner | Where it is described |
|---|---|---|
| Is this user still registered, and under this uid? | Admission::admit, called by the connector |
Admission |
| Tell the connection whose flows it carries | the connector, through the LeaseSlot |
this page |
| Pick an outbound | Router::route with a RouteTarget built by route_target |
this page, and Route model |
| Refuse audited destinations and record the hit | RuleManager::detect, through Dispatcher::forbidden |
this page |
| Dial the outbound | Outbound::connect_stream and Outbound::connect_datagram |
Inbounds and outbounds |
| Bill bytes and hold them back while the user is in debt | Gate and Metered |
Metering |
| Route, audit and bill each UDP packet | FanOut |
this page |
| Queue effects, poll links, report failures to the core | the server runtime in etemenanki-concepts |
Server runtime |
The connector spawns no task and owns no socket of its own. Everything it returns, including the FanOut and every sub-link inside it, is polled by the one task that runs the connection’s runtime. Dropping the runtime’s outbound slot drops all of it.
Key types
Section titled “Key types”The kernel contract
Section titled “The kernel contract”The connector implements the kernel’s Connector trait for the Flow type every server core in etemenanki-protocols produces.
pub trait Connector<Target> { type Stream: AsyncRead + AsyncWrite + Unpin; type Datagram: DatagramLink; type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;}
pub trait DatagramLink: Unpin { type Addr;
fn poll_send_to( &mut self, cx: &mut Context<'_>, buf: &[u8], to: &Self::Addr, ) -> Poll<io::Result<usize>>;
fn poll_recv_from( &mut self, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>, ) -> Poll<io::Result<Self::Addr>>;}
pub enum Outbound<S, D> { Stream(S), Datagram(D),}pub struct Flow<T> { pub destination: Destination, pub user: NetworkUser<T>, pub sniffed: Option<SniffedBehavior>, pub source: Option<IpAddr>,}For a UDP flow, destination is the first packet’s. Later packets carry their own address, which reaches the datagram link as the to argument of poll_send_to. sniffed is set only when the inbound sniffs and the destination named an IP; see Sniffing.
Dispatcher
Section titled “Dispatcher”One Dispatcher is shared by every connection of one listener. TransportManager::start in src/manager/transport.rs builds it next to the listener, and its router is never swapped: a route edit replaces the node’s transport, and its connections with it (route_change_drops_connections in tests/unit/e2e.rs).
pub struct Dispatcher { pub router: Arc<Router<Outbound>>, pub rules: Arc<RuleManager>, pub node_tag: CompactString, pub admission: Admission,}
impl Dispatcher { fn forbidden(&self, dest: &Destination, uid: i64) -> bool;}
fn dest_string(d: &Destination) -> String;forbidden calls RuleManager::detect(&self.node_tag, &dest_string(dest), Some(uid)). dest_string yields the destination host as the audit rules match it: the domain name, or the IP literal, without the port.
pub fn detect(&self, tag: &str, dest: &str, uid: Option<i64>) -> bool;detect takes the rule map’s RwLock for reading, tries the node’s rules in order, and on the first regex that matches inserts (uid, rule_id) into the node’s HashSet of hits under a Mutex. Because the ledger is a set, recording the same hit again, for example when a client keeps sending UDP packets to the same forbidden host, adds nothing. The reporting loop takes the set with RuleManager::drain.
KatanaConnector and LeaseSlot
Section titled “KatanaConnector and LeaseSlot”pub type LeaseSlot = watch::Sender<Option<CancellationToken>>;
#[derive(Clone)]pub struct KatanaConnector { disp: Arc<Dispatcher>, source: Option<IpAddr>, lease: Option<Arc<LeaseSlot>>,}
impl KatanaConnector { pub fn new( disp: Arc<Dispatcher>, source: Option<IpAddr>, lease: Option<Arc<LeaseSlot>>, ) -> Self;}
type ConnectFuture = Pin< Box<dyn Future<Output = io::Result<link::Outbound<Metered<OutboundStream>, FanOut>>> + Send>,>;
impl Connector<Flow<UserTag>> for KatanaConnector { type Stream = Metered<OutboundStream>; type Datagram = FanOut; type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;}A connector is cheap to clone and there is one per connection. It is built in two places:
| Call site | source |
lease |
How a retired user’s connection ends |
|---|---|---|---|
serve_stream in src/serve.rs (every TCP-based inbound) |
the peer address the listener saw | Some: the sender of a watch channel whose receiver the connection’s driver watches |
the driver waits on the published lease and ends the whole connection when it is cancelled |
run_hysteria in src/serve.rs |
Some(ip), the client’s address |
None |
the listener owns its runtimes, so each flow stops when its Gate refuses to move bytes |
Outbounds and refusal
Section titled “Outbounds and refusal”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, dest: &Destination) -> StreamFuture; pub fn connect_datagram(&self, dest: &Destination) -> DatagramFuture;}
pub fn refused() -> io::Error;refused() is io::Error::new(io::ErrorKind::PermissionDenied, "refused"). The message deliberately gives no reason, and every refusal in this file uses it: an unknown user, a blocked route and an audit hit look the same from outside.
OutboundDatagram in src/outbound/proxy.rs boxes whatever datagram link an outbound opened:
pub struct OutboundDatagram(Box<dyn DatagramLink<Addr = Destination> + Send>);
impl OutboundDatagram { pub fn new<D: DatagramLink<Addr = Destination> + Send + 'static>(link: D) -> Self;}Gate and Metered
Section titled “Gate and Metered”Both are covered in Metering. The connector uses them like this:
impl Gate { pub fn new(counter: Arc<UserCounter>, retired: CancellationToken) -> Self; pub fn poll_open(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>; pub fn sent(&self, n: usize); pub fn received(&self, n: usize);}
pub struct Metered<S> { inner: S, gate: Gate,}
impl<S> Metered<S> { pub fn new(inner: S, gate: Gate) -> Self;}poll_open is ready once the user’s lease is not cancelled and their token bucket is out of debt. A cancelled lease makes it fail with ConnectionAborted and the message the user was retired. sent and received add to the user’s counter and charge the bucket.
Connecting a flow
Section titled “Connecting a flow”Everything up to the dial happens synchronously inside connect, before the returned future is first polled. Only the dial itself is asynchronous.
flowchart TB
A["connect(flow)"] --> B{"admission.admit(tag)"}
B -- None --> R1["Err(refused())"]
B -- "Some(counter, lease)" --> C["publish the lease once"]
C --> D["source = flow.source or the listener's peer"]
D --> E["Gate::new(counter, lease)"]
E --> F{"destination.network"}
F -- Udp --> G["Datagram(FanOut::new(...))"]
F -- "any other" --> H["router.route(route_target(dest, sniffed, source))"]
H --> I{"Block, or forbidden by audit?"}
I -- yes --> R2["Err(refused())"]
I -- no --> J["outbound.connect_stream(dest).await"]
J --> K["Stream(Metered::new(stream, gate))"]
-
Admit.
connectclones theArc<UserTag>inflow.user.user_dataand callsself.disp.admission.admit(&tag). It returns the user’sArc<UserCounter>and lease only if the credential is still registered to that uid. OnNonethe connector returnsErr(refused())without logging: a departed user’s client keeps retrying until it notices, and a line per attempt would flood the log. -
Publish the lease once. When the connector has a
LeaseSlot, it callssend_if_modifiedwith a closure that stores the lease only if the slot is stillNone, and reports a change only then. The first admitted flow therefore tells the connection which user it belongs to; later flows neither replace the lease nor wake the receiver. This happens before routing, so a connection whose first flow is then refused still learns its lease. -
Pick the source address.
flow.source.or(self.source): a source the inbound decoded for this flow wins over the peer the listener saw. -
Build the gate.
Gate::new(counter, lease)binds the user’s counter and token bucket to the lease. Each flow has its ownGate; flows of one user share the counter and bucket behind it. -
UDP: hand over a
FanOut. Whenflow.destination.network == DialNetwork::Udp, the connector returnslink::Outbound::Datagram(FanOut::new(disp, gate, tag.uid, source))at once. Nothing is routed, audited or dialed yet: the association has no single destination, so all of that happens per packet. -
TCP: route. Any other network takes the stream path.
route_target(&flow.destination, flow.sniffed.as_ref(), source)builds the target, andRouter::routereturns the first matching rule’s outbound or the default. -
Refuse. If the outbound is
Outbound::Block, orforbidden(&flow.destination, tag.uid)is true, the connector returnsErr(refused()). The check short-circuits: a flow the router blocks is never matched against the audit rules, so a blocked destination records no audit hit. -
Dial and meter.
outbound.connect_stream(&flow.destination)starts the dial. The returned future awaits it, propagates its error with?, and wraps the stream aslink::Outbound::Stream(Metered::new(stream, gate)).
The route target
Section titled “The route target”pub fn route_target<'a>( dest: &'a Destination, sniffed: Option<&'a SniffedBehavior>, source: Option<IpAddr>,) -> routing::RouteTarget<'a>;impl<Outbound> Router<Outbound> { pub fn route(&self, target: RouteTarget<'_>) -> Arc<Outbound>;}route_target converts the destination’s Remote and port, and then adds what it has:
| Input | Builder call | Why |
|---|---|---|
sniffed: Some(s) |
with_sniffed_domain(s.domain.as_str()) |
A flow addressed by IP carries its real name only in its payload; this is what lets domain and geosite rules reach it. |
source: Some(ip) |
with_source(ip) |
Feeds RouteMatch::SourceCidr, when the listener knows the client. |
dest.network is Tcp or Udp |
with_network(TargetNetwork::Tcp) or with_network(TargetNetwork::Udp) |
Feeds RouteMatch::Network; a matcher whose field is None never matches. Unknown and Unix add nothing. |
The two callers pass different inputs:
| Caller | dest |
sniffed |
source |
|---|---|---|---|
KatanaConnector::connect (stream path) |
flow.destination |
flow.sniffed.as_ref() |
the flow’s source, else the peer |
FanOut::poll_send_to (each packet) |
the packet’s to |
None |
the same source, fixed for the association |
Packets are never sniffed, so a UDP packet addressed by IP can only match IP-based rules (cidr, geoip) and port rules.
UDP fan-out
Section titled “UDP fan-out”A UDP association can talk to many peers, and each packet names its own. FanOut is the datagram link katana hands the runtime for it. It routes, audits, gates and bills each packet on its own, and keeps one sub-link open per outbound the association has used.
pub const MAX_SUBS: usize = 64;
struct Sub { key: usize, link: OutboundDatagram,}
pub struct FanOut { disp: Arc<Dispatcher>, gate: Gate, uid: i64, source: Option<IpAddr>, subs: Vec<Sub>, next: usize, opening: Option<(usize, DatagramFuture)>, recv_waker: Option<Waker>,}
impl FanOut { fn new(disp: Arc<Dispatcher>, gate: Gate, uid: i64, source: Option<IpAddr>) -> Self; fn key_of(outbound: &Arc<Outbound>) -> usize; fn poll_opening(&mut self, cx: &mut Context<'_>) -> Poll<()>;}
impl DatagramLink for FanOut { type Addr = Destination; // poll_send_to, poll_recv_from}| Field | Meaning |
|---|---|
disp |
The generation’s dispatcher: router, audit rules and node tag. |
gate |
The user’s speed limit and retirement, checked before every send and receive. |
uid |
The uid audit hits are recorded against. |
source |
The client address used in every packet’s route target. |
subs |
Open sub-links, least recently sent to at the front and most recently sent to at the back. |
next |
The index the next receive starts at, so every sub gets its turn. |
opening |
The one sub-link being opened, with its key. |
recv_waker |
The receiver’s waker, parked while no sub can deliver. |
A sub-link’s key is Arc::as_ptr(outbound) as usize. The router holds every outbound Arc for the whole generation, and the FanOut holds the router through disp, so the pointer is a stable name for “the same outbound” for as long as the association lives.
Sending a packet
Section titled “Sending a packet”flowchart TB
S["poll_send_to(buf, to)"] --> R["route(route_target(to, None, source))"]
R --> X{"Block, or forbidden?"}
X -- yes --> D1["Ready(Ok(len)): dropped, not billed"]
X -- no --> G{"gate.poll_open"}
G -- Pending --> P["Pending: the runtime keeps the packet queued"]
G -- Err --> E["Err(ConnectionAborted)"]
G -- Ok --> L{"sub with this key?"}
L -- yes --> T["sub.poll_send_to, then on Ok(n): gate.sent(n) and move the sub to the back"]
L -- no --> O{"opening"}
O -- None --> N["opening = connect_datagram(to)"]
N --> Q{"poll_opening, this key"}
O -- "this key" --> Q
O -- "another key" --> Q2{"poll_opening, other key"}
Q -- Pending --> P
Q -- opened --> L
Q -- failed --> D2["Ready(Ok(len)): dropped, not billed"]
Q2 -- Pending --> P
Q2 -- "opened or failed" --> L
In words, poll_send_to:
- Routes the packet with
route_target(to, None, self.source). - Drops it if the outbound is
Blockorforbidden(to, self.uid)is true, by returningPoll::Ready(Ok(buf.len())). UDP has no refusal to give one peer, and the association carries others, so failing the send would be wrong. The gate is not consulted and nothing is billed. As on the stream path, a blocked packet is not matched against the audit rules. - Waits on
gate.poll_open(cx). While the user is in debt the call isPendingand the runtime keeps the packet at the head of its queue. Every poll starts from the top, so a packet that waited is routed and audited again when the runtime retries it. The audit result is recorded only on a hit, and the ledger is a set, so a client that keeps sending to the same forbidden host adds a single(uid, rule_id)entry until the node’s reporting loop drains the set withRuleManager::drain, not one per packet. - Looks for a sub with the outbound’s key. If there is one, it sends through it. On
Ready(Ok(n))it callsgate.sent(n)and moves the sub to the back ofsubs. Any other result,Pendingor an error, is returned as is. - With no matching sub, it drives
opening:None: it startsoutbound.connect_datagram(to)and loops.- The same key: it polls the open.
Pendingstays pending. A success adds the sub and the loop sends through it. A failure drops the packet withOk(buf.len()), unbilled, so the association does not stall on an outbound that cannot be reached. - Another key: the packet polls that open until it finishes, successfully or not, and then loops to start its own. Only one open is ever in flight.
poll_recv_from also drives the open in flight, as described under Receiving below. When the receive side finishes an open that failed, the waiting packet finds opening empty and no sub for its key on its next poll, so it starts a new open instead of being dropped.
Opening, eviction and removal
Section titled “Opening, eviction and removal”poll_opening polls the in-flight DatagramFuture. On success it pushes the new Sub to the back of subs. If that makes more than MAX_SUBS, it removes subs[0], the least recently sent-to sub, and decrements next with saturating_sub(1) so the round-robin position keeps pointing at the same sub; when next pointed at the evicted sub, it stays 0 and points at the new front. It then wakes recv_waker, because the new sub has never been polled for a reply and so has no waker registered with its socket. On failure it logs udp fan-out: opening an outbound failed: {e} at debug level and clears opening.
stateDiagram-v2 [*] --> Opening: packet routed to an outbound with no sub Opening --> Open: connect_datagram succeeds, pushed to the back Opening --> [*]: connect_datagram fails, packet dropped Open --> Open: a send succeeds, moved to the back Open --> [*]: evicted when a new sub makes MAX_SUBS + 1 Open --> [*]: poll_recv_from fails, removed
A failed open is not remembered. The next packet routed to the same outbound tries again. Evicting a sub drops its link: its socket or tunnel closes, and replies still addressed to it are lost.
What each outbound gives connect_datagram:
Outbound variant |
Sub-link | Packets are addressed |
|---|---|---|
Direct |
FreedomConnector::bind_udp, a ResolvingUdp; dest is not used |
per packet |
Socks |
SocksUdpLink, UDP ASSOCIATE over a new control stream to the upstream; dest is not used |
per packet |
Wireguard |
the WgConnector’s datagram link for an anonymous flow |
per packet |
Vmess, Vless |
a proxy client runtime dialed toward dest, the packet that opened it |
to the opening destination |
Http |
Err, Unsupported: http carries no datagrams |
none |
ShadowsocksLegacy, Shadowsocks2022 |
Err, Unsupported: shadowsocks carries no datagrams, shadowsocks-2022 carries no datagrams |
none |
Block |
not reached: FanOut drops the packet before opening anything |
none |
With the HTTP and Shadowsocks outbounds, every packet routed to them costs one failed open attempt and is dropped.
Receiving
Section titled “Receiving”poll_recv_from:
- Waits on
gate.poll_open(cx). While the user is in debt, no sub is read, and replies wait in the sub-links’ sockets or tunnels. - Calls
poll_opening(cx), so the in-flight open also makes progress when only the receive side is being polled. - Polls each sub once, starting at
nextand wrapping around. The first sub that delivers wins:nextmoves to the sub after it,gate.received(n)bills the bytes it wrote intobuf, and its sender address is returned. - If a sub fails, it logs
udp fan-out: a sub-link ended: {e}at debug level, removes that sub, resetsnextto0if it now points past the end, and wakes itself withcx.waker().wake_by_ref()so the remaining subs are polled again. The error is not returned: the association outlives any one of its sub-links. - Stores the waker in
recv_wakerand returnsPending.
Starting at next instead of 0 keeps one busy sub from starving the others.
An association end to end
Section titled “An association end to end”sequenceDiagram participant RT as Server runtime participant KC as KatanaConnector participant FO as FanOut participant OB as Outbound participant SL as Sub-link RT->>KC: connect(flow with network Udp) KC-->>RT: Datagram(FanOut) RT->>FO: poll_send_to(packet, to) FO->>FO: route, audit, gate.poll_open FO->>OB: connect_datagram(to) OB-->>FO: OutboundDatagram FO->>SL: poll_send_to(packet, to) SL-->>FO: Ok(n) FO->>FO: gate.sent(n) RT->>FO: poll_recv_from(buf) FO->>SL: poll_recv_from, round robin from next SL-->>FO: reply, from FO->>FO: gate.received(n) FO-->>RT: Ok(from)
What a packet costs
Section titled “What a packet costs”| What happens to the packet | Billed | Audit hit recorded |
|---|---|---|
Routed to Block |
no | no |
| Matches an audit rule | no | yes, (uid, rule_id) once |
| The open of its sub-link fails | no | no |
| Accepted by a sub-link | yes, gate.sent(n) |
no |
Accepted and then dropped inside the sub-link, for example a name with no usable address in ResolvingUdp |
yes: the sub-link reports it as sent | no |
| A reply read from any sub-link | yes, gate.received(n) |
no |
The direct outbound’s UDP socket
Section titled “The direct outbound’s UDP socket”Outbound::Direct wraps a FreedomConnector. Its UDP side is ResolvingUdp, a dual-stack socket that resolves domain targets itself, so a single direct sub-link serves every destination the association names.
const MAX_RESOLVED_NAMES: usize = 256;
impl FreedomConnector { pub fn new(resolver: Resolver, strategy: AddressFamilyStrategy) -> Self; pub async fn connect(&self, dest: &Destination) -> io::Result<TcpStream>; pub fn bind_udp(&self) -> io::Result<ResolvingUdp>;}
type ResolveFuture = Pin<Box<dyn Future<Output = Option<IpAddr>> + Send>>;
pub struct ResolvingUdp { socket: DualStackUdp, support: FamilySupport, resolver: Resolver, strategy: AddressFamilyStrategy, resolved: HashMap<CompactString, Option<IpAddr>>, order: VecDeque<CompactString>, resolving: Option<(CompactString, ResolveFuture)>,}
impl ResolvingUdp { fn poll_target(&mut self, cx: &mut Context<'_>, to: &Destination) -> Poll<Option<SocketAddr>>;}Binding
Section titled “Binding”bind_udp asks the dialer’s bind_dual for one socket per family the outbound’s address_family allows:
AddressFamilyStrategy |
Families bound |
|---|---|
Ipv4Only |
V4 |
Ipv6Only |
V6 |
Auto, PreferIpv4, PreferIpv6 |
V4 and V6 |
bind_dual in environment/src/dial/udp.rs logs a family that fails to bind at debug level and carries on with the other; it fails only when no family bound, with AddrNotAvailable and udp: no usable local socket in any requested family. bind_udp then records which sockets actually bound in a FamilySupport. Only those families can send, so resolution later skips addresses of a family with no socket and takes a usable address further down the answer instead of dropping the packet. An error from bind_dual fails the open, and FanOut drops the packet.
Choosing a packet’s address
Section titled “Choosing a packet’s address”poll_target returns Ready(Some(addr)), Ready(None) for “drop it”, or Pending:
- An IP literal is used as is when
strategy.allows(ip)andsupport.supports(ip). Otherwise it isNone. - A name already in
resolvedreturns its cached answer, which may beNone. - A new name starts a lookup if none is in flight:
resolve_candidates("freedom", &dest, strategy, support, &resolver), keeping the first candidate. The candidates are already filtered by policy and bound families, and ordered by theprefer_*strategy. A lookup error becomesNone. - Another name’s lookup in flight makes the packet wait:
poll_targetreturnsPendingwithout polling that lookup. Only one lookup runs at a time per sub-link. The lookup advances only when a packet for its own name is polled, which the runtime does by retrying the packet at the head of the key’s effect queue.
When a lookup completes, the answer is inserted into resolved and the name is pushed onto order. When order already holds MAX_RESOLVED_NAMES names, the oldest inserted name is removed first. Eviction is first in, first out: a cache hit does not refresh a name’s position.
The cache has no expiry of its own. A name keeps the answer it first got, including a negative one, until 256 later names push it out or the sub-link closes. The shared Resolver underneath has its own TTL-bound cache; see DNS.
poll_send_to sends to the chosen address through the dual-stack socket. When poll_target returns None, it logs freedom: dropping a datagram with no usable address at debug level and returns Ok(buf.len()). poll_recv_from reads from either socket and reports the sender as Destination::udp(addr).
Invariants
Section titled “Invariants”| Invariant | Mechanism | Pinned by |
|---|---|---|
| A credential the registry does not hold, or holds under another uid, opens nothing | Admission::admit looks up (key, uid) before anything else; None returns refused() |
a_user_the_registry_does_not_know_is_refused, a_credential_rebound_to_another_uid_is_refused in tests/unit/connector.rs |
| The connection learns its user from the first admitted flow, and that lease is cancelled when the user leaves | LeaseSlot::send_if_modified stores only into an empty slot; Admission::commit cancels the lease |
the_lease_reaches_the_connection_and_goes_with_the_user in tests/unit/connector.rs |
| A retired user’s connection ends | the driver in src/serve.rs awaits the published lease |
a_retired_users_connection_ends in tests/unit/serve.rs |
| A retired user moves no more bytes, even on a parked read | Gate::poll_open polls the lease’s cancelled_owned future first |
a_retired_users_stream_refuses_to_move, retiring_the_user_wakes_a_parked_read in tests/unit/meter.rs |
| A stream is billed to its user in both directions | Metered calls Gate::sent and Gate::received after each completed transfer |
an_admitted_stream_is_billed_to_its_user in tests/unit/connector.rs; each_direction_is_billed_to_the_user in tests/unit/meter.rs |
| The speed limit does not depend on transfer size | the bucket is charged the full byte count after each transfer and blocks further transfers while in debt | the_limit_holds_however_the_writes_are_sized, the_limit_is_shared_by_both_directions in tests/unit/meter.rs |
A TCP flow the router blocks is refused with PermissionDenied |
matches!(*outbound, Outbound::Block) before the dial |
a_blocked_destination_is_refused in tests/unit/connector.rs |
| A TCP flow to an audited destination is refused, and the hit is recorded against the user | Dispatcher::forbidden → RuleManager::detect |
a_forbidden_destination_is_refused_and_recorded in tests/unit/connector.rs; detect_records_and_drains in tests/unit/rule.rs |
| Each UDP packet is routed on its own; a blocked packet is dropped and costs nothing; a delivered one is billed both ways | per-packet route in FanOut::poll_send_to; gate.sent only after a sub-link accepts |
udp_is_billed_after_routing_and_blocked_packets_are_free in tests/unit/connector.rs |
| A UDP packet to an audited destination is dropped unbilled, and the hit is recorded | per-packet forbidden before the gate |
a_forbidden_udp_destination_is_dropped_and_recorded in tests/unit/connector.rs |
| A sniffed name reaches domain rules for a TCP flow addressed by IP | route_target adds with_sniffed_domain |
a_sniffed_host_reaches_a_domain_rule, disable_sniffing_stops_the_domain_rule_matching in tests/integration/sniff.rs |
An association keeps at most MAX_SUBS sub-links |
poll_opening removes subs[0] when a push exceeds the cap |
no dedicated test |
| An association has at most one open in flight | opening is a single Option; another key waits for it |
no dedicated test |
A direct sub-link remembers at most MAX_RESOLVED_NAMES names and runs one lookup at a time |
order eviction and the single resolving slot in ResolvingUdp::poll_target |
no dedicated test |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Where | What happens | What the runtime does with it |
|---|---|---|
admit returns None |
Err(refused()), no log line |
sends the core Event::ConnectFailed |
TCP route is Block, or audit forbids |
Err(refused()) |
sends the core Event::ConnectFailed |
connect_stream fails |
the dial’s io::Error, returned by the connect future |
sends the core Event::ConnectFailed |
| UDP packet blocked or forbidden | Ready(Ok(buf.len())), dropped |
counts it as sent |
| Sub-link open fails | debug log, Ready(Ok(buf.len())), dropped |
counts it as sent |
| Sub-link send fails | the error is returned from poll_send_to |
drops the packet, keeps the association, sends the core Event::SendFailed |
| Sub-link receive fails | debug log, sub removed, Pending |
nothing: the association continues on the other subs |
| User retired, on send | Err(ConnectionAborted), the user was retired |
drops the packet and sends the core Event::SendFailed |
| User retired, on receive | Err(ConnectionAborted) |
fails the whole outbound key and sends the core Event::OutboundError |
| Direct: no usable address for a packet | debug log, Ready(Ok(buf.len())) |
counts it as sent |
Cancellation needs no extra code. FanOut owns its sub-links, the in-flight DatagramFuture and the Gate; ResolvingUdp owns its sockets and the in-flight ResolveFuture. When the runtime drops the outbound slot, because the core closed the key, the connection ended, or the generation was torn down, dropping them cancels the pending dial and lookup and closes every socket and tunnel.
Limits
Section titled “Limits”| Name | Value | Where | Effect |
|---|---|---|---|
MAX_SUBS |
64 |
src/connector.rs |
Sub-links one UDP association keeps open. The least recently sent-to sub is closed when a new one opens past the cap. |
| Opens in flight | 1 per association | FanOut::opening |
A packet for another outbound waits until the current open finishes. |
MAX_RESOLVED_NAMES |
256 |
src/outbound/freedom.rs |
Names a direct UDP sub-link remembers. A client naming more pays another lookup for the evicted ones. |
| Lookups in flight | 1 per direct sub-link | ResolvingUdp::resolving |
A packet to another new name waits until the current lookup finishes. |
Unit tests for this file live in tests/unit/connector.rs, included into src/connector.rs as its tests module. They build a real Dispatcher over a pool with a direct and a block outbound, one registered user (uid 1), and an optional CIDR rule to block. They exercise real loopback TCP and UDP echo servers, so billing is checked against bytes that crossed a socket.
| Test | What it pins |
|---|---|
a_user_the_registry_does_not_know_is_refused |
An unregistered credential gets PermissionDenied. |
a_credential_rebound_to_another_uid_is_refused |
A credential now registered to another uid gets PermissionDenied, so the new account is never billed for the old one’s flows. |
an_admitted_stream_is_billed_to_its_user |
A TCP flow opens as a Metered stream, echoes, and bills 5 bytes up and 5 down. |
a_blocked_destination_is_refused |
A destination routed to block gets PermissionDenied. |
a_forbidden_destination_is_refused_and_recorded |
An audit regex hit refuses the flow and drains as DetectResult { uid: 1, rule_id: 9 }. |
the_lease_reaches_the_connection_and_goes_with_the_user |
The first flow publishes the lease; committing an empty user set cancels it; the next flow is refused. |
udp_is_billed_after_routing_and_blocked_packets_are_free |
A UDP flow opens as a FanOut; a delivered packet and its echo bill 4 bytes each way; a packet to a blocked CIDR returns its length and adds nothing. |
a_forbidden_udp_destination_is_dropped_and_recorded |
The association still opens for a forbidden first destination; the packet is not billed; the hit drains as DetectResult { uid: 1, rule_id: 4 }. |
Related tests in other files: tests/unit/meter.rs for Gate and Metered, tests/unit/rule.rs for RuleManager (detect_records_and_drains, update_skips_identical_ruleset, unknown_uid_records_negative_one), and tests/unit/serve.rs for a_retired_users_connection_ends. The sniffing tests in tests/integration/sniff.rs run a real Xray client against a katana node and return early when no Xray binary is available.
MAX_SUBS eviction, the single in-flight open and ResolvingUdp have no dedicated tests. A change to any of them should add one, with a stub outbound whose connect_datagram can be held pending or made to fail.