The plane: routing each flow
Source files: 39 · checked against Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/router.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/build/route.rsEtemenanki/supervisor/src/build/dns.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/topology/spec_plan/route.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/policy.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/topology/outbound/mod.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/topology/outbound/dns.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/build/validate.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/environment/src/routing.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/protocols/src/socks/server.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/concepts/src/link.rsEtemenanki/concepts/src/runtime.rsEtemenanki/supervisor/tests/unit/plane.rsEtemenanki/supervisor/tests/unit/router.rsEtemenanki/supervisor/tests/unit/plan.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/app/tests/integration/e2e_route_context.rsEtemenanki/app/tests/integration/e2e_sniff.rsEtemenanki/app/tests/support/mod.rskatana/src/rule.rskatana/src/manager/node.rs
Every flow the supervisor carries is routed on the plane: one immutable snapshot of the compiled route table, the outbound and balancer targets its rules name, and the DNS service that port-53 flows are handed to. The snapshot sits behind a single ArcSwap cell. An apply that changes routing, outbounds or DNS builds a new snapshot and publishes it with one pointer store; every connector reads the cell again for each new flow. This is read-copy-update, and it is why a route change never has to stop a listener or touch a connection that is already open.
This page is for contributors who change supervisor/src/connector.rs, supervisor/src/topology/{plane,router,flow}.rs or supervisor/src/build/route.rs. It covers the plane and its cell, how a route table is compiled into slots, the types a flow carries (Flow, Principal, FlowContext), AppConnector::connect step by step, targets and the Guarded wrapper that lets a drain end a flow, and epochs. How an apply decides to publish, and the drain policies, are on Planning and applying a change; the matchers themselves are on Route model; the outbounds a target wraps, the UDP fan-out and balancers are on Outbounds, UDP fan-out and balancers.
Responsibilities
Section titled “Responsibilities”The plane and the connector:
- hold one consistent view of routes, targets and the DNS service, replaced as a whole (
PlaneCell,Plane); - admit each flow on its session before anything else happens (
Session::admit); - route each TCP flow once, on the plane current when
connectis called, and a UDP association packet by packet, on the plane current at eachFanOutLink::poll_send_tocall; - intercept port 53 for the supervisor’s own DNS service, when there is one, over TCP and UDP alike;
- resolve a balancer to one member per flow, so the flow is filed under, and drained with, the outbound that carries it;
- open the flow on its target and wrap it so a drain can end it (
Target::connect_stream,Guarded); - register the opened flow with the tracker, with its rule, outbound version and plane epoch (
Tracker::meter_stream).
They do not:
- evaluate matchers.
CompiledRoutes::decidehands aRouteTargettoetemenanki_environment::routing::Router, which does the first-match walk (Route model). - dial anything themselves. The
Outboundbehind a target dials (Outbounds, UDP fan-out and balancers). - decide when a plane is published or which versions drain. The planner and the actor’s commit do (Planning and applying a change).
- sniff. The protocol cores, the mux demuxer and the SOCKS driver fill
Flow::sniffedbefore they open the flow (Sniffing). - log.
connectand the plane write no log line, and a TCP flow’s failures return to the caller asio::Errors. The UDP fan-out logs its sub-link failures at debug level (Log lines, and Outbounds, UDP fan-out and balancers).
katana relies on the rule that a flow is filed under its resolved outbound. Its destination audit turns each panel rule into a route toward a blackhole outbound tagged audit#<id>, subscribes to the tracker, and records a hit for each FlowEvent::Opened it receives whose outbound().tag parses as audit#<id> (src/rule.rs → audit_hit, src/manager/node.rs → spawn_audit). The audit rules match domain names only, requested or sniffed.
Who reads the plane
Section titled “Who reads the plane”| Reader | Loads the cell | Routes with | For |
|---|---|---|---|
AppConnector::connect (supervisor/src/connector.rs) |
once per call, for every flow that is not UDP (in practice TCP) | Plane::route |
every TCP flow opened through a connector: a connection’s single flow, each mux sub-flow, each Hysteria 2 proxy stream, each SOCKS CONNECT. A UDP flow passes through connect once, to become a fan-out, without loading the cell. |
FanOutLink::poll_send_to (supervisor/src/topology/outbound/udp_fanout.rs) |
on every call | Plane::route |
the packets of every UDP association |
RoutedDialer::dial (supervisor/src/topology/outbound/dns.rs) |
once per connection to an upstream DNS server | Plane::route_upstream |
the split resolver’s through-proxy queries |
Supervisor::epoch (supervisor/src/supervisor.rs) |
once per call | none | reporting the current plane’s epoch |
A runtime calls AppConnector::connect synchronously while it applies Effect::Open (concepts/src/runtime.rs; The server runtime). A mux carrier’s core pushes one Open per sub-flow, so every sub-flow goes through connect and a TCP sub-flow loads the plane again. The SOCKS driver (protocols/src/socks/server.rs) runs no runtime: it calls connect and awaits the future itself, once for a CONNECT, and once per UDP ASSOCIATE, on the first datagram of the association that it forwards (SOCKS).
Key types
Section titled “Key types”PlaneCell and Plane
Section titled “PlaneCell and Plane”pub const DNS_PORT: u16 = 53;
pub type PlaneCell = Arc<ArcSwap<Plane>>;
pub struct Plane { routes: Arc<CompiledRoutes>, /// The target of each route slot, in `CompiledRoutes::targets` order. targets: Box<[Arc<Target>]>, /// The service applications' DNS is answered by, when the supervisor /// answers it itself. dns: Option<Arc<Target>>, epoch: u64,}
impl Plane { pub(crate) fn new( routes: Arc<CompiledRoutes>, targets: Box<[Arc<Target>]>, dns: Option<Arc<Target>>, epoch: u64, ) -> Self; pub(crate) fn empty() -> Self; pub fn epoch(&self) -> u64; pub fn route(&self, flow: &Flow, ctx: &FlowContext) -> Routed<'_>; pub fn route_upstream(&self, flow: &Flow, ctx: &FlowContext) -> Routed<'_>;}PlaneCell is an Arc<arc_swap::ArcSwap<Plane>> (arc-swap 1.9). One cell exists per supervisor: Actor::new creates it with ArcSwap::from_pointee(Plane::empty()) and puts it in Shared, which every accept loop clones; the Supervisor handle keeps a clone for epoch. Readers call load(), which returns a guard without taking a lock; the actor calls store() to publish. A Plane is never mutated after it is built.
targetsis indexed by slot.Plane::newchecks withdebug_assert_eq!that it has exactly one entry per slot ofroutes.dnsisSomeonly when the DNS spec isDnsSpec::Split, the one mode in which the supervisor answers applications’ queries itself (Name resolution and the DNS service).Plane::empty()is the value the cell holds before the first commit: a table with one slot, the empty tag, filled byOutbound::Blackholeunder the id@v0, no DNS target, epoch 0.the_empty_plane_drops_everythingpins that every flow, port 53 included, is dropped by it.Supervisor::startapplies the first spec before it returns, and that commit stores the first real plane before it starts any listener, so no connection of the supervisor’s own ever reads the empty plane.
route, route_upstream and Routed
Section titled “route, route_upstream and Routed”#[derive(Clone, Copy)]pub struct Routed<'a> { pub target: &'a Arc<Target>, /// The rule that matched, if one did. pub rule: Option<RuleId>,}#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]pub struct RuleId(u32);
impl RuleId { pub const fn new(raw: u32) -> Self; pub const fn get(self) -> u32;}route is what applications’ flows go through:
- If the plane has a DNS target and
flow.destination.port == DNS_PORT, it returns that target withrule: None. The check reads the port alone: not the address, and not the network. The code comment gives the reason: an application retries over TCP when a UDP answer comes back truncated, and the retry must reach the same service. - Otherwise it calls
route_upstream.
route_upstream asks the table for a Decision and returns &self.targets[decision.slot] with decision.rule. It never intercepts. The DNS service’s own upstream connections are routed with it, because intercepting them would send the service’s questions back to itself (port_53_is_intercepted_except_for_the_services_own_queries).
A RuleId is the matching rule’s index in RouteSpec::rules of the plane that routed the flow. After an apply that changes the rules, the same number can name another rule; the flow’s plane epoch says which plane’s rules it indexes.
Routed borrows its target from the plane, so it lives no longer than the guard load() returned. Every caller clones the Arc<Target> before the guard is dropped: connect through Target::resolve, the fan-out with routed.target.clone() before it resolves, and RoutedDialer with .target.clone().
Target and TargetKind
Section titled “Target and TargetKind”pub struct Target { id: OutboundId, kind: TargetKind, /// Cancelled to close the flows opened on this version, when its drain /// policy says so. closed: CancellationToken,}
pub enum TargetKind { Outbound(Outbound), /// Chooses among outbound targets by health, per flow. Balancer(Arc<Balancer>),}
impl Target { pub fn outbound(id: OutboundId, outbound: Outbound) -> Self; pub fn balancer(id: OutboundId, balancer: Arc<Balancer>) -> Self; pub fn id(&self) -> &OutboundId; pub fn kind(&self) -> &TargetKind; pub fn resolve(self: &Arc<Self>, prefer: Option<&OutboundId>) -> Arc<Target>; pub fn connect_stream(self: &Arc<Self>, flow: Flow) -> impl Future<Output = io::Result<Guarded<OutboundStream>>> + Send + 'static; pub fn connect_datagram(self: &Arc<Self>, flow: Flow) -> impl Future<Output = io::Result<Guarded<OutboundDatagram>>> + Send + 'static; pub(crate) fn close_flows(&self);}
fn nested_balancer(id: &OutboundId) -> io::Error;A target is one version of one tag: an outbound built from its spec, or a balancer over outbound targets. Each has its own CancellationToken, created with the target, so the flows opened on version 2 of direct can be closed without touching those on version 3. TargetKind carries #[allow(clippy::large_enum_variant)]: a target lives behind an Arc for its whole life, so boxing the larger variant would only add an indirection.
| Method | What it does |
|---|---|
resolve(prefer) |
An outbound target returns itself (an Arc clone). A balancer returns the member target Balancer::select_preferring(prefer) picks now. prefer is the member a caller already uses for this route; under round robin the balancer keeps it while it is healthy, so a UDP association does not hop between members packet by packet. The connector passes None. |
connect_stream(flow) |
Resolves with None when it is called, not when its future is first polled, so a balancer picks the member at the call. Then it calls Outbound::connect_stream(flow) on the member or outbound, takes a clone of that target’s closed token, and returns a future that wraps the opened stream in Guarded. |
connect_datagram(flow) |
The same for a UDP flow, over Outbound::connect_datagram. The fan-out resolves a balancer itself (it keys sub-links by the member’s version) and then calls this on the member, where resolve is a no-op. |
close_flows() |
Cancels closed. Every Guarded holding a clone fails on its next poll. Only the actor’s drain calls it. |
Selection happens per flow rather than at route time, so a member its health probe has marked down is skipped from the next flow on. How members are probed, how they are chosen, and what a balancer does when every member is down are on Outbounds, UDP fan-out and balancers.
If a resolved target were itself a balancer, both connect methods would fail with balancer member <tag>@v<n> is itself a balancer (nested_balancer, an io::Error::other). Validation makes that unreachable: every balancer member must be an outbound tag with a probeable upstream.
Guarded
Section titled “Guarded”pub struct Guarded<T> { inner: T, closed: Pin<Box<WaitForCancellationFutureOwned>>,}
impl<T> Guarded<T> { fn new(inner: T, closed: CancellationToken) -> Self; fn poll_closed(&mut self, cx: &mut Context<'_>) -> Poll<io::Error>;}Guarded is how a drain reaches a flow that is already relaying. Guarded::new holds the inner link and closed.cancelled_owned(), boxed and pinned, built from the clone of the target’s token. poll_closed is the one check the guarded methods share; each of these fails once the token is cancelled:
| Trait | Guarded methods |
|---|---|
AsyncRead |
poll_read |
AsyncWrite |
poll_write, poll_flush |
DatagramLink (Addr = Destination) |
poll_send_to, poll_recv_from |
Once the token is cancelled, each of them returns io::ErrorKind::ConnectionAborted with the text the outbound this flow was opened on was drained. A read or receive that is parked on a silent upstream is woken by the cancel and fails at once (closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive). poll_shutdown is passed straight through without the check.
For a stream under a runtime, the runtime treats the error like any other outbound error: the core receives Event::OutboundError for that key and closes that flow. A single-flow core finishes the connection; a mux carrier closes the one sub-flow (The server runtime). A SOCKS relay gets the error from the stream it relays. A guarded UDP sub-link’s error stays inside its FanOutLink, which drops that sub-link (Outbounds, UDP fan-out and balancers).
CompiledRoutes and Decision
Section titled “CompiledRoutes and Decision”pub struct CompiledRoutes { table: routing::Router<Decision>, /// The tag each slot names, indexed by slot. targets: Vec<CompactString>,}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]pub struct Decision { /// The slot the flow is routed to. pub slot: usize, /// The rule that matched, or `None` for the default route. pub rule: Option<RuleId>,}
impl CompiledRoutes { pub fn new( rules: impl IntoIterator<Item = (Vec<routing::RouteMatch>, CompactString)>, default: CompactString, geo: GeoData, ) -> Self; pub fn single(tag: CompactString) -> Self; pub fn targets(&self) -> &[CompactString]; pub fn decide(&self, flow: &Flow, ctx: &FlowContext) -> Decision;}The compiled table names slots, not outbounds. Each distinct tag that a rule or the default names gets one slot, and the plane fills the slots with the targets of the same publish. Keeping outbounds out of the table is what lets an outbound be rebuilt without recompiling the table, and lets one compiled table be reused while outbounds change: the plane is rebuilt around the same Arc<CompiledRoutes> with new slot contents.
CompiledRoutes::new walks the rules in order:
slot_of(tag)looks the tag up in aHashMap<CompactString, usize>and, the first time a tag is seen, appends it totargetsand returns its index. Slots are therefore numbered in order of first mention, and the default’s tag is added last if no rule named it.- Each rule becomes one
RouteItemwhose output is its ownArc<Decision>with the rule’s slot andrule: Some(RuleId::new(index)). Rules that share a tag share the slot but keep separate decisions, so a flow always says which rule sent it. The index is converted withu32::try_from(at).expect("fewer than 2^32 rules"). - The default is
Arc<Decision>with its slot andrule: None. - The table is
routing::Router::new(RouteTable { routes, default }, geo).
single(tag) is new([], tag, GeoData::default()), the table Plane::empty and the unit tests use. decide builds the RouteTarget with route_target, calls Router::route, and copies the Decision out of the Arc it returns.
For the rules [port 80 → proxy, network udp → direct, port 443 → proxy] with default direct, every_rule_naming_a_tag_shares_that_tags_slot pins:
| Flow | targets() |
Decision |
|---|---|---|
| TCP to port 80 | ["proxy", "direct"] |
slot of proxy, rule 0 |
| TCP to port 443 | slot of proxy, rule 2 |
|
| UDP to port 443 | slot of direct, rule 1 (the first matching rule wins) |
|
| TCP to port 22 | slot of direct, None (the default) |
route_target
Section titled “route_target”pub fn route_target<'a>(flow: &'a Flow, ctx: &'a FlowContext) -> routing::RouteTarget<'a>;route_target adapts a supervisor flow to the config-agnostic target the route model matches against. It borrows everything; nothing is allocated.
RouteTarget field |
Taken from | Matchers that read it |
|---|---|---|
remote |
flow.destination.remote, as Remote::Domain(&str) or Remote::IpAddr |
domain, GeoSite, Cidr, GeoIp |
port |
flow.destination.port |
PortRange |
sniffed_domain |
flow.sniffed’s domain, when set |
every domain matcher and GeoSite, as an alternative to remote |
inbound_tag |
ctx.inbound_tag, always |
InboundTag |
source |
flow.source.or(ctx.source) |
SourceCidr |
network |
TargetNetwork::Tcp or Udp from flow.destination.network; not set for DialNetwork::Unknown or Unix |
Network |
Two choices here carry reasons in the code:
- The flow’s own source wins over the listener’s. The code comment gives the reason for that order: a listener that serves many peers from one socket, such as a QUIC listener, knows the peer only per flow. On every accept path at this revision the two carry the same address.
serve_connectionhandsAppConnector::source()to the protocol core or the SOCKS driver, which copies it onto each flow it opens (Flow::new(dest, user, source)); the Hysteria 2 listener builds one connector per QUIC connection and the TUN listener one per TCP flow or UDP association, each with the same IP their cores put on the flows. Only a Unix socket supplies neither, andRoutedDialer’s flows have none. The fourroute_targettests insupervisor/tests/unit/router.rspin every combination. - The network is always supplied. Withholding it would leave every
RouteMatch::Networkrule permanently inert.
A matcher whose field is absent does not match. A flow on a Unix-socket inbound has no source, so no SourceCidr rule matches it.
Flow, Principal and FlowContext
Section titled “Flow, Principal and FlowContext”pub type Flow = etemenanki_protocols::flow::Flow<Principal>;
pub struct Principal { user: Option<UserKey>, label: CompactString, revoked: OnceLock<UserRemovalPolicy>,}
impl Principal { pub fn user(user: UserKey, label: CompactString) -> Arc<Self>; pub fn anonymous() -> Arc<Self>; pub fn user_key(&self) -> Option<UserKey>; pub fn label(&self) -> &str; pub fn revoked(&self) -> Option<UserRemovalPolicy>; pub(crate) fn revoke(&self, policy: UserRemovalPolicy);}
pub fn anonymous_user() -> NetworkUser<Principal>;
#[derive(Clone, Debug)]pub struct FlowContext { pub inbound_tag: CompactString, pub source: Option<IpAddr>,}Flow is the protocols crate’s flow with Principal as its user payload:
pub struct Flow<T> { pub destination: Destination, pub user: NetworkUser<T>, /// What sniffing recovered from the flow's first bytes, when the inbound /// sniffs and the destination named only an IP. pub sniffed: Option<SniffedBehavior>, /// The client's address, when the inbound knows it. pub source: Option<IpAddr>,}
impl<T> Flow<T> { pub fn new(destination: Destination, user: NetworkUser<T>, source: Option<IpAddr>) -> Self; pub fn toward(&self, destination: Destination) -> Self;}NetworkUser<Principal> carries the authorization the client presented and user_data: Arc<Principal>. Flow::new is the constructor the protocol cores, the TUN inbound, the SOCKS driver and RoutedDialer use; it starts with sniffed: None. For a UDP flow, destination is the first packet’s address; later packets carry their own. Flow::toward(destination) makes a flow toward another destination with the same user and source and sniffed: None; the fan-out uses it to route each packet as its own flow, and the mux demuxer to open each sub-flow. Clone is written out so that cloning never requires Principal: Clone.
A Principal is who a flow belongs to: the user an inbound admitted it as, or nobody. Every server protocol’s user table hands one back on a successful handshake, so a session learns its user from the first flow it opens. Each inbound holds its own principal per user, so revoking a user on one inbound (revoke, first call wins, stored in a OnceLock) does not touch that user’s sessions elsewhere. The connector uses the principal twice: Session::admit binds the session to it on the first flow, refusing a revoked one, and FlowMeta keeps it so the tracker can find the user’s speed limit and label. The admission rules and the revocation race are on Users, principals and sessions.
Principal::anonymous() has no user key and an empty label. It serves inbounds in their open or shared-credential mode and flows the supervisor opens itself. anonymous_user() wraps it in a NetworkUser whose authorization is UsernamePassword with an empty username and password; RoutedDialer and the unit tests use it.
FlowContext is where a flow entered. Both fields exist only because a rule can match on them: carrying them keeps RouteMatch::InboundTag and RouteMatch::SourceCidr from being rules that can never fire. source is the peer the listener captured at accept, when it has one.
AppConnector and FlowScope
Section titled “AppConnector and FlowScope”#[derive(Clone)]pub struct AppConnector { plane: PlaneCell, ctx: FlowContext, tracker: Tracker, /// The session the flows belong to. A TUN inbound's flows have none. session: Option<Arc<Session>>, /// Whether datagram flows count toward the session's wire bytes. charge_datagrams: bool,}
impl AppConnector { pub(crate) fn new( plane: PlaneCell, ctx: FlowContext, tracker: Tracker, session: Option<Arc<Session>>, ) -> Self; pub(crate) fn charging_datagrams(self) -> Self; pub fn source(&self) -> Option<IpAddr>; pub fn session(&self) -> Option<&Arc<Session>>;}
// `Outbound` here is `etemenanki_concepts::link::Outbound`.type ConnectFuture = Pin< Box< dyn Future<Output = io::Result<Outbound<Metered<Guarded<OutboundStream>>, FanOutLink>>> + Send, >,>;
impl Connector<Flow> for AppConnector { type Stream = Metered<Guarded<OutboundStream>>; type Datagram = FanOutLink; type Future = ConnectFuture;
fn connect(&mut self, flow: Flow) -> ConnectFuture;}
#[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>>,}AppConnector implements the Connector<Flow> trait from concepts/src/link.rs (fn connect(&mut self, target) -> Self::Future, resolving to link::Outbound::Stream or link::Outbound::Datagram). Only the crate constructs it. A front end never builds one; it gets the supervisor’s connectors by running inbounds. Inside the crate, serve_connection reads source() and hands it to the protocol core or the SOCKS driver as the client address; session() is read by serve_connection’s SOCKS branch, which wraps the stream in WireCounted to count its bytes toward the session’s wire, and by drive.
charge_datagrams is set, through charging_datagrams(), only for SOCKS: a SOCKS UDP association’s datagrams never cross the session’s own stream, so their payload is counted toward the session’s wire instead. Every other inbound’s datagrams travel inside the session’s transport and are already counted there. Per-user usage accounting explains the wire.
Where connectors come from (Listeners and the serve loop has the accept paths):
| Accept path | One connector per | ctx.source |
session |
|---|---|---|---|
Stream listener over TCP (Connection::serve_socket) |
session: the accepted socket’s, or over gRPC each HTTP/2 stream’s, which is a session of its own | the peer’s IP | that session |
| Stream listener over a Unix socket | accepted socket | None |
the socket’s session |
Hysteria 2 (run_hysteria_inbound) |
QUIC connection | the client’s IP | the connection’s session. A connection that finishes its handshake after the listener was stopped gets None and an already-cancelled token, so it is closed at once and no session escapes the close of a removed inbound. |
TUN (run_tun_inbound) |
TCP flow, and UDP association per source address and port | the host’s IP | None: a device admits no users |
Stream and Hysteria 2 connectors are built through Shared (supervisor/src/serve.rs). Shared::connector registers a new session with Sessions::open and builds its connector; Sessions::open returns None once the listener’s stop token is cancelled, and then no connector is built. Shared::connector_for builds a connector for a session already registered, such as the one the stream accept loop opens for each accepted socket. TUN connectors, and the Hysteria 2 connector of a connection that arrives after the stop, are built with AppConnector::new and no session.
Data flow
Section titled “Data flow”The shape of a plane
Section titled “The shape of a plane”flowchart LR cell["PlaneCell (ArcSwap)"] --> plane["Plane, epoch N"] plane --> routes["Arc CompiledRoutes"] plane --> slots["targets, one per slot"] plane --> dnst["dns: Option Target"] routes -- "Decision.slot" --> slots slots --> direct["Target direct@v2"] slots --> lb["Target lb@v1 (balancer)"] lb --> a["member a@v1"] lb --> b["member b@v3"] dnst --> svc["Target @vN, Outbound::Dns"]
A balancer target holds its members’ targets through its Members: the same Arc<Target>s the actor holds for those outbound tags, and that a slot holds when a rule also names a member directly. The DNS target is internal: its tag is empty and its version is the epoch of the plane it was first published in.
Routing one flow
Section titled “Routing one flow”flowchart TB
start["Plane::route(flow, ctx)"] --> dns{"DNS target and port 53?"}
dns -- yes --> svc["DNS target, rule None"]
dns -- no --> up["route_upstream"]
up --> rt["route_target(flow, ctx)"]
rt --> pick["Router::route, first match"]
pick -- "rule i matched" --> d1["Decision: slot, Some(i)"]
pick -- "no rule" --> d0["Decision: default slot, None"]
d1 --> slot["targets[slot]"]
d0 --> slot
slot --> res["Target::resolve"]
res --> member["outbound target (itself, or a balancer member)"]
One TCP flow through connect
Section titled “One TCP flow through connect”sequenceDiagram participant RT as Runtime or SOCKS driver participant AC as AppConnector participant S as Session participant P as PlaneCell participant T as Target participant TR as Tracker RT->>AC: connect(flow) AC->>S: admit(principal) AC->>P: load() P-->>AC: guard on the current Plane AC->>AC: route, resolve(None), FlowMeta::of AC->>T: connect_stream(flow) T-->>AC: future (outbound dial, closed token) AC-->>RT: boxed future, guard dropped RT->>AC: poll the future T-->>AC: Guarded stream AC->>TR: meter_stream(meta, stream) AC-->>RT: link::Outbound::Stream(Metered)
AppConnector::connect, step by step
Section titled “AppConnector::connect, step by step”- Admit. When the connector has a session,
session.admit(&flow.user.user_data)runs first, for TCP and UDP alike. A refusal is returned as an already-completed future, so the caller sees a failed connect. Users, principals and sessions has the rules; the refusal texts are under Failure paths and cancellation. - UDP becomes a fan-out. A flow whose
destination.networkisDialNetwork::Udpis not routed here. The connector buildsFanOutLink::new(plane, ctx, flow, FlowScope { tracker, session: session.id(), charge }), wherechargeis the session’s wire only whencharge_datagramsis set, and returnslink::Outbound::Datagram(link)at once. No flow is registered yet; each sub-link registers when it opens. - Load the plane. Every other flow, whatever its network, goes on:
let plane = self.plane.load();The comment in the code states the rule this step implements: loaded once for this flow, so a route change reaches the next flow while this one keeps the target it is opened on. - Route.
plane.route(&flow, &self.ctx)returns the target and the matching rule, or the DNS target when port 53 is intercepted. - Resolve.
routed.target.resolve(None). A balancer picks its member now, so the flow is filed under the outbound that carries it. - Describe.
FlowMeta::of(&flow, session id, &ctx.inbound_tag, ctx.source, target.id().clone(), routed.rule, plane.epoch())fixes what the tracker will report: session, inbound tag, principal, source (flow.source.or(ctx.source)), requested destination, sniffed domain, outbound version, rule and plane epoch. - Open.
target.connect_stream(flow)calls the outbound’sconnect_streamand captures a clone of the target’sclosedtoken. The plane guard is dropped whenconnectreturns; the boxed future carries the outbound’s own future, the token clone, aTrackerclone and theFlowMeta. - Meter. When the caller polls the future to completion, the opened stream is wrapped as
Metered<Guarded<OutboundStream>>bytracker.meter_stream(meta, stream), which registers the flow and emitsFlowEvent::Opened. A dial that fails returns its error through?before this step, so a failed open is never registered as a flow.
Steps 1 to 7 run synchronously inside the caller’s task, in the connect call: a runtime’s call that applies Effect::Open, or the SOCKS driver’s. connect spawns nothing, and step 8 runs when the caller polls the returned future.
UDP: routed packet by packet
Section titled “UDP: routed packet by packet”FanOutLink::poll_send_to routes each time it is called: self.flow.toward(to), self.plane.load(), plane.route(...), then resolve(prefer) with the member the association already uses for that route, and a sub-link keyed by the resolved OutboundId. Sub-links are keyed by the id rather than by the address of an Arc<Target>, because outbounds outlive a single apply, one tag can have an old version draining beside a new one, and a freed allocation can be reused by its successor (the module comments of supervisor/src/entity/id.rs and udp_fanout.rs). A packet to port 53 is intercepted exactly as a TCP flow is, because the fan-out calls the same route.
Each sub-link is a tracked flow of its own, with a FlowMeta built from Flow::toward(to):
FlowEntry::destination()is the destination of the packet that opened the sub-link;FlowEntry::sniffed()is alwaysNone, becausetowardclears it;FlowEntry::rule()is the rule that routed that packet, orNonefor the default route or a query the supervisor’s DNS service answers;FlowEntry::plane_epoch()is the epoch of the plane that routed that packet.
a_udp_association_counts_each_sub_link_and_charges_its_session pins that two destinations routed to two outbounds give two flows, one with rule() None and one with rule 0. The sub-link table, its size (MAX_SUBS = 64) and its send and receive rules are on Outbounds, UDP fan-out and balancers.
DNS interception
Section titled “DNS interception”With DnsSpec::Split, build_dns returns a Service, and prepare wraps it as Target::outbound(internal_id(epoch + 1), Outbound::Dns(service)). internal_id builds an OutboundId with the empty tag; validation refuses an empty tag in any spec, so the id never collides with a spec outbound. The version is the epoch the same commit is about to store, so each rebuilt DNS target has a version no earlier one had.
| Flow | Route result | Opened as |
|---|---|---|
| TCP to any address, port 53 | DNS target, rule: None |
OutboundStream::Proxy over a DnsTcpStream on the service, ready at once with no dial |
| UDP packet to any address, port 53 | DNS target, rule: None |
a DnsUdpLink sub-link on the service, ready at once with no dial |
Any flow, port 53, with DnsSpec::Single |
the route table, like any port | whatever the table names |
The resolver’s own connection to an upstream server (RoutedDialer) |
route_upstream: the route table only |
whatever the table names for the server’s address |
RoutedDialer is built only for a split resolver with a non-empty through_proxy list, which uses the direct resolver as its bootstrap (supervisor/src/build/dns.rs). It routes each upstream connection as a TCP flow toward the server’s IP and port, with the anonymous user and FlowContext { inbound_tag: DNS_FLOW_TAG, source: None }, where DNS_FLOW_TAG is "dns", so a rule can match the resolver’s traffic by the inbound tag dns. It opens the connection with Target::connect_stream, so the stream is Guarded and a closing drain reaches it; it never goes through the tracker, so these flows are not registered, killable or charged to anyone. It holds the cell as Weak<ArcSwap<Plane>>, because the plane owns, through its outbounds, the resolver that owns the dialer. dial upgrades the reference synchronously; once nothing holds the plane cell any more (the supervisor, its accept loops and every live connector and fan-out link), the upgrade fails and the returned future fails with dns: the supervisor this resolver dials through is gone when it is polled. The service itself is on Name resolution and the DNS service and DNS resolver.
Sniffed and requested destinations
Section titled “Sniffed and requested destinations”Each opener that sniffs sets flow.sniffed just before it opens the flow:
| Opener | Where the sniffed name comes from |
|---|---|
| A sniffing protocol core (HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022, the Hysteria 2 stream core) | its sniff prefix, set right before the core pushes Effect::Open, in its open_sniffed or inline just before it opens |
PassthroughCore, which serves TUN flows (protocols/src/core/mod.rs) |
its sniff prefix, in open |
The mux demuxer (protocols/src/mux/demux.rs) |
the payload of the sub-flow’s New frame, through crate::sniff::sniff |
The SOCKS driver (protocols/src/socks/server.rs) |
the prefix it collects after granting the request, before it calls connect |
The plane uses the sniffed name only as an extra input to domain matchers (RouteTarget::sniffed_domain). Nothing on the connect path rewrites flow.destination: the outbound dials the address the client requested, and no outbound reads sniffed. FlowMeta keeps both, so FlowEntry::destination() reports the requested address and FlowEntry::sniffed() the name.
The app’s end-to-end tests make this observable. In app/tests/integration/e2e_sniff.rs the only route to freedom is a domain_suffix rule for sniffed.example, the default is blackhole, the client asks for 127.0.0.1, and sniffed.example does not resolve. A successful echo therefore proves both that the sniffed name matched and that the flow was dialed to the requested IP. Sniffing covers the sniffers.
Compiling the route table
Section titled “Compiling the route table”pub(crate) fn compile_routes(spec: &RouteSpec) -> io::Result<CompiledRoutes>;pub struct RouteSpec { pub rules: Vec<RouteRuleSpec>, pub default: CompactString, pub geoip: Option<PathBuf>, pub geosite: Option<PathBuf>,}
pub struct RouteRuleSpec { pub matchers: Vec<RouteMatch>, pub outbound: CompactString,}A rule matches when any one of its matchers does, rules are tried in order, and a flow no rule matches takes default. Both outbound and default name an outbound or a balancer: the two share one tag namespace. geoip and geosite name the files compile_routes reads through build_geo_data.
compile_routes runs in Actor::prepare for the plan’s Step::Build(Resource::Route). check, the dry run behind etemenanki-app --test, plans from a fresh actor, so it always compiles. It:
-
collects every
RouteMatch::GeoSite(code)andRouteMatch::GeoIp(code)across all rules, in rule order, duplicates included; -
calls
build_geo_data(spec.geoip, spec.geosite, &geosite_codes, &geoip_codes)(environment/src/routing.rs), which:- handles geosite before geoip, so a spec whose geosite and geoip references both fail reports the geosite error;
- reads a geo file only when at least one matcher of its kind is used, so a configured path that no matcher uses is never opened (
etemenanki-app --testprintsConfiguration OK.for ageoippath that does not exist when only a port rule is configured); - loads each distinct matcher string once, skipping one already in the map, and keys the result by that raw string, for example
google@adsor!cn; - looks a code up ignoring ASCII case: the part before
@for geosite, and the part after a leading!for geoip;
Route model describes the decoding and what
@attributeand!do; -
hands the rules as
(matchers, outbound)pairs, the default and the geo data toCompiledRoutes::new.
Any error fails the apply as ApplyError::Build { resource: Resource::Route, .. } and leaves the running plane as it was. Validation runs before this, so a rule or default naming a tag that is neither an outbound nor a balancer never reaches it. As printed by etemenanki-app --test, which prefixes each with configuration invalid::
| Cause | Text |
|---|---|
| A rule or the default names an unknown tag (validation) | route references unknown outbound <tag> |
| A geosite matcher with no geosite file | building route failed: a geosite matcher is used but no geosite file is configured |
| A geoip matcher with no geoip file | building route failed: a geoip matcher is used but no geoip file is configured |
| A code missing from the file | building route failed: geosite code not found: <code> or building route failed: geoip code not found: <code>. <code> is the name without its @attribute or leading !: nosuchcode@ads and !nosuchcode both print nosuchcode. |
| A file that is not a geo list | building route failed: geosite decode: <decoder error> (or geoip decode: …) |
| A file that cannot be read | building route failed: <OS error>, for example No such file or directory (os error 2) |
| A geosite code whose regex entries still fail to compile together after the invalid ones are skipped | building route failed: geosite <code>: regex entries: <error> |
| A geoip entry whose address is not 4 or 16 bytes | building route failed: geoip cidr: ip must be 4 or 16 bytes, got <n> |
A geoip entry whose prefix does not fit in a u8 |
building route failed: geoip cidr: prefix too large |
| A geoip entry whose prefix is too long for its address | building route failed: geoip cidr: <error> |
Two geosite problems are logged at warn (target etemenanki_environment::routing) and skipped rather than refused: a domain entry of an unknown type (skipping unknown geosite domain type <n>) and a regex entry that does not compile (skipping invalid geosite regex "<pattern>": <error>).
Invalid port ranges and domain regexes are refused earlier, while the front end lowers its rules with parse_port_match and parse_domain_regexes (etemenanki-app: from TOML to a spec).
Publishing a plane
Section titled “Publishing a plane”The planner adds Step::PublishPlane when the DNS spec, an outbound (added, changed, removed, or rebuilt by a DNS change), a balancer (added, changed, removed, or rebuilt because one of its members was) or the route spec changed. User sets and inbounds never publish, and user edits through set_users, upsert_user and remove_user never go through the planner or touch the cell. Planning and applying a change has each case and the plan tests that pin it.
In Actor::commit, the user keys and speed limits are committed and self.targets is replaced by the targets prepare built or reused. Then, only when the plan contains Step::PublishPlane:
self.epoch += 1.- Fill the slots: for each tag in
routes.targets(), in slot order, takeself.targets[tag].clone(). Every such tag exists, because validation refused any rule naming a missing one. self.shared.plane.store(Arc::new(Plane::new(routes.clone(), slots, dns_target.clone(), self.epoch))). From this store on, everyload()sees the new plane.- Cancel the probe token of every balancer the plan does not
Reuse. For each new balancer, create a child of the root token and callBalancer::spawn_probeon the supervisor’sTaskTracker, resolving with the newdns.serversresolver and connecting withTcpDialer::new(socket)under the supervisor’s socket policy (Outbounds, UDP fan-out and balancers).
On every commit, with or without a publish, the actor then stores routes, dns and dns_target as its running ones. Without PublishPlane, the cell, the epoch and the probes are left alone.
Everything else in the commit runs after the store, including starting new listeners, and the plan’s Drain steps come last. Two consequences follow. A listener started by an apply never serves a connection on the previous plane. An old outbound version is drained only once new flows already go to its successor; a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish asserts that PublishPlane comes before the Drain step. The commit ends with the debug line applied: <n> built, <n> reused, <n> swapped, <n> drained.
A reused outbound or balancer is the same Arc<Target> in the old plane and the new one, with the same token. That is how an unchanged stateful outbound, such as a Hysteria 2 client with its QUIC connection, carries across an apply (an_unchanged_hysteria2_outbound_keeps_its_quic_connection_across_apply in supervisor/tests/hot_swap.rs).
Epochs
Section titled “Epochs”| Where | Meaning |
|---|---|
Plane::epoch() |
The epoch this plane was published under: 0 for Plane::empty(), 1 for the plane of the first spec, one more for each later publish. |
Supervisor::epoch() |
self.plane.load().epoch(): the current plane’s epoch. It reads the cell directly and does not go through the actor. |
FlowEntry::plane_epoch() |
The epoch of the plane that routed the flow, from FlowMeta. For a UDP sub-link, the plane that routed the packet that opened it. |
| DNS target version | The epoch of the plane the target was first published in. |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow asserts that changing the default route moves Supervisor::epoch() to exactly one more; a_refused_spec_changes_nothing asserts that a refused spec leaves it unchanged. Because every tracked flow carries its epoch, a front end can tell flows routed before an apply from those routed after it, for example to end every flow an older plane routed:
let current = supervisor.epoch();let killed = supervisor .tracker() .kill_where(|flow| flow.plane_epoch() < current);What open flows keep
Section titled “What open flows keep”A route change never moves an open stream: it stays on the target it was opened on. An open stream carries the link its outbound opened and, inside Guarded, a clone of that target’s closed token; the token is what lets a later drain reach it. The stream never reads the cell again, and the load() guard its connect took was released when connect returned. A UDP association holds the cell itself and reads it on every send.
| Change after a flow opened | TCP flow, or one mux sub-flow | UDP association |
|---|---|---|
| A route change sends its destination elsewhere | Stays on the target it was opened on (a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow). The next sub-flow of the same mux carrier takes the new route (a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session). |
The next packet is routed on the new plane and goes where the new route says (the same UDP test). |
| Its outbound is rebuilt or removed | The actor applies the old version’s drain policy to that version’s token (Drains). | The actor applies the drain policy to the existing sub-link’s token in the same way; a packet whose route now names the new version opens a sub-link there, because the ids differ. |
| Its balancer is rebuilt | Nothing: the flow is guarded by its member’s token, not the balancer’s (a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers). |
Later packets resolve on the new balancer. |
| A balancer member goes down | Nothing: selection happens when a flow opens. | Later packets may resolve to another member. |
| A user’s speed limit changes | Felt at the flow’s next read or write (Tracking flows, stats and speed limits). | The same, per sub-link. |
Drains
Section titled “Drains”An apply that rebuilds or removes an outbound plans Step::Drain { outbound, policy } for the old version, after PublishPlane; DrainPolicy::Keep is the default. In commit, drain (supervisor/src/supervisor.rs) runs for the old target when the actor’s previous targets map holds that exact version. The policy decides only what the actor does with that version’s closed token:
DrainPolicy |
The version’s token |
|---|---|
Keep |
Never cancelled by the actor, so Guarded never fails these flows. The new plane no longer names the version, so no new flow is opened on it. |
Close |
Cancelled at once by target.close_flows(). |
CloseAfter(grace) |
Cancelled by close_flows() after grace, from a task on the supervisor’s TaskTracker that selects on tokio::time::sleep(grace) and the root token; the task ends without cancelling when the root token is cancelled first. |
How the policy is resolved for each outbound, and what the apply report lists, are on Planning and applying a change.
a_changed_outbound_keeps_its_old_flows_unless_drained_with_close exercises the first two with a freedom outbound: under the default the flow on version 1 still echoes after the apply; after a second change under Close, the flow on version 2 is closed, the flow on version 1 is not drained again, and a new flow uses version 3. shutdown_ends_pending_grace_timers pins that a CloseAfter of an hour neither outlives nor delays shutdown.
Only outbound versions are drained. A balancer target’s own token is never cancelled by the supervisor, and it does not need to be: a flow through a balancer is guarded by its member’s token, so a Close or CloseAfter drain of that member’s version reaches it. The DNS target is never drained either; no Drain step names it.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| Routes and targets are published together, and a rule never names a tag the plane lacks. | One store of a Plane built from one apply’s routes and targets; validation refuses unknown route tags; debug_assert_eq! on the slot count. |
a_refused_spec_changes_nothing (an unknown default is refused and the epoch is unchanged); every_rule_naming_a_tag_shares_that_tags_slot |
| One slot per tag, and a decision still names its rule. | slot_of in CompiledRoutes::new; one Arc<Decision> per rule. |
every_rule_naming_a_tag_shares_that_tags_slot |
The first matching rule wins; no match takes the default with rule: None. |
RouteTable::pick. |
every_rule_naming_a_tag_shares_that_tags_slot; the route model’s own tests |
| A TCP flow is routed once; a route change reaches the next flow, not an established one. | connect loads the cell once and the opened stream never consults it again. |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
| Each mux sub-flow is routed when it opens. | The runtime calls connect per Effect::Open. |
a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session |
| A UDP association is routed on the plane current at each send. | FanOutLink::poll_send_to loads the cell on every call. |
a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow |
| Port 53 reaches the DNS service over TCP and UDP when one exists, and the service’s own queries never do. | Plane::route versus Plane::route_upstream. |
port_53_is_intercepted_except_for_the_services_own_queries, without_a_dns_service_port_53_follows_the_route_table |
| Before the first commit, nothing leaves the host. | Plane::empty. |
the_empty_plane_drops_everything |
| Inbound tag, source and network reach the router. | route_target. |
inbound_tag_selects_the_route, source_cidr_matches_the_client_address, network_separates_tcp_from_udp (app/tests/integration/e2e_route_context.rs); the four source tests in supervisor/tests/unit/router.rs |
| A sniffed name is matched, and the requested destination is dialed. | route_target sets sniffed_domain; nothing rewrites flow.destination. |
an_http_host_routes_an_ip_addressed_flow, a_tls_sni_routes_an_ip_addressed_flow, a_flow_whose_sniffed_host_does_not_match_is_blocked, an_unsniffable_payload_falls_through_to_the_default, turning_sniffing_off_stops_the_domain_rule_matching (app/tests/integration/e2e_sniff.rs) |
| A balanced flow is filed under, and drained with, the member that carries it. | resolve(None) before FlowMeta::of; connect_stream takes the resolved target’s token. |
a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers |
| A closing drain fails every direction of a flow and wakes a parked receive. | Guarded::poll_closed in every guarded method but poll_shutdown. |
closing_a_targets_flows_aborts_its_open_streams, closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive |
| An old version is drained only after the new plane is published. | The planner orders Drain after PublishPlane; commit walks the steps after the store. |
a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish |
| A flow is admitted before it is routed. | Session::admit is the first statement of connect. |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy (supervisor/tests/unit/session.rs) covers the admission itself |
| A flow is registered only once its outbound opened, with its rule, resolved outbound and epoch. | meter_stream runs after opened.await?. |
Registration after success: no dedicated test. The metadata: mux_sub_flows_and_their_session_count_exactly_what_moved, a_udp_association_counts_each_sub_link_and_charges_its_session (supervisor/tests/tracking.rs) |
Versions of a tag never repeat, so keying by OutboundId never confuses a draining build with its successor. |
next_version; versions of removed tags are kept in RunningState. |
each_rebuild_of_a_tag_takes_the_next_version, a_tag_added_back_takes_a_version_it_never_had |
| A refused apply leaves the plane as it was. | Routes and targets are built in prepare; only commit stores. | a_refused_spec_changes_nothing |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Where | Error | What follows |
|---|---|---|
Session::admit: the session’s token is cancelled, or it is gone from the registry |
PermissionDenied, the session was closed |
The future resolves to the error at once. A runtime hands its core Event::ConnectFailed for the key; the SOCKS driver gets the error from its await (SOCKS). |
Session::admit: the first flow presents a revoked principal |
PermissionDenied, the user was removed before the session opened a flow |
The session’s token is cancelled as well, which ends the whole session. |
| A resolved target is a balancer | balancer member <tag>@v<n> is itself a balancer |
Unreachable while validation holds. |
| The outbound’s dial fails | The outbound’s error | Returned before meter_stream; no flow is registered. The texts are on Outbounds, UDP fan-out and balancers. |
| A drain closes the flow’s version | ConnectionAborted, the outbound this flow was opened on was drained |
Returned by the flow’s next read, write, flush, send or receive. Under a runtime the core gets Event::OutboundError and closes that flow. A UDP sub-link that fails this way is dropped from its association. |
| A tracker kill | ConnectionAborted, the flow was killed |
Raised by Metered, outside Guarded (Tracking flows, stats and speed limits). |
RoutedDialer once nothing holds the plane cell |
NotConnected, dns: the supervisor this resolver dials through is gone |
The resolver’s query fails. |
compile_routes fails in prepare |
ApplyError::Build for route |
The apply is refused; the plane is not republished. |
| A rule names an unknown tag | ApplyError::UnknownReference, route references unknown outbound <tag> |
Refused by validation, before prepare. |
Log lines
Section titled “Log lines”| Where | Level and target | Text |
|---|---|---|
FanOutLink, a sub-link fails to open |
debug, etemenanki_supervisor::topology::outbound::udp_fanout |
udp fan-out: opening <tag>@v<n> failed: <error> |
FanOutLink, a send on a sub-link fails |
debug, the same |
udp fan-out: sending on <tag>@v<n> failed: <error> |
FanOutLink, a receive on a sub-link fails |
debug, the same |
udp fan-out: sub-link <tag>@v<n> ended: <error> |
| A balancer probe sees a member’s health change | info, etemenanki_supervisor::topology::balancer |
balancer member <tag> is now up or balancer member <tag> is now down |
Actor::commit, at its end |
debug, etemenanki_supervisor::supervisor |
applied: <n> built, <n> reused, <n> swapped, <n> drained |
build_geo_data |
warn, etemenanki_environment::routing |
skipping unknown geosite domain type <n>, skipping invalid geosite regex "<pattern>": <error> |
connect, the plane, Guarded and RoutedDialer log nothing.
Cancellation
Section titled “Cancellation”connect and the plane spawn nothing. The caller owns the returned future (a runtime keeps it in the key’s LinkState::Connecting), and dropping it cancels the dial.
The commit that publishes a plane starts the balancer probe tasks, one per member of each new balancer, and a CloseAfter drain starts a timer task. Both run on the supervisor’s TaskTracker. A probe task ends when its balancer’s probe token, a child of the root token, is cancelled: at the publish that no longer reuses the balancer, or at shutdown. A CloseAfter task ends after its grace or when the root token is cancelled. The probe loop is on Outbounds, UDP fan-out and balancers.
Limits
Section titled “Limits”| Item | Where | Value |
|---|---|---|
| Intercepted port | supervisor/src/topology/plane.rs → DNS_PORT |
53, over TCP and UDP |
| Rules per table | supervisor/src/topology/router.rs → CompiledRoutes::new |
fewer than 2^32; RuleId is a u32 and compiling more panics with fewer than 2^32 rules |
| Matchers stored inline per rule | environment/src/routing.rs → RouteItem |
3 (SmallVec<[RouteMatch; 3]>); more are stored on the heap, with no cap |
| Cell reads | supervisor/src/connector.rs, supervisor/src/topology/outbound/udp_fanout.rs, supervisor/src/topology/outbound/dns.rs |
one ArcSwap::load per TCP flow (a mux TCP sub-flow included), per poll_send_to call, and per resolver connection |
| Epochs and versions | supervisor/src/supervisor.rs, supervisor/src/topology/spec_plan/plan.rs |
u64 counters, one step per publish and per rebuild of a tag |
| Sub-links per UDP association | supervisor/src/topology/outbound/udp_fanout.rs → MAX_SUBS |
64 (Outbounds, UDP fan-out and balancers) |
A routing decision costs one walk of the rules in order until the first match, with each flow’s domains lower-cased once per decision; the route model’s page gives the per-matcher costs.
| Layer | File | What it covers |
|---|---|---|
| Unit, plane | supervisor/tests/unit/plane.rs |
port_53_is_intercepted_except_for_the_services_own_queries, without_a_dns_service_port_53_follows_the_route_table, the_empty_plane_drops_everything, closing_a_targets_flows_aborts_its_open_streams, closing_a_targets_flows_aborts_its_datagram_links_and_wakes_a_waiting_receive, a_balanced_stream_is_closed_by_its_members_drain_not_the_balancers. Targets are blackholes, so no network is involved. |
| Unit, router | supervisor/tests/unit/router.rs |
the_circuits_own_source_wins_over_the_listeners, without_one_the_listeners_source_is_used, with_neither_there_is_no_source_to_match_on, a_circuit_source_works_with_no_listener_source, every_rule_naming_a_tag_shares_that_tags_slot. |
| Unit, plan | supervisor/tests/unit/plan.rs |
When PublishPlane is planned, for example the_first_plan_binds_and_builds_everything_then_publishes, an_unchanged_spec_reuses_everything_and_publishes_nothing and removing_a_balancer_alone_republishes_the_plane; that Drain follows it (a_changed_outbound_is_a_new_version_and_the_old_one_drains_after_the_publish); version numbering (each_rebuild_of_a_tag_takes_the_next_version, a_tag_added_back_takes_a_version_it_never_had). Planning and applying a change lists the rest. |
| Hot swap | supervisor/tests/hot_swap.rs |
Real sockets reconfigured while connections are live: a_route_change_reaches_the_next_udp_packet_but_not_an_established_tcp_flow, a_route_change_reaches_the_next_mux_sub_flow_of_a_live_session, a_changed_outbound_keeps_its_old_flows_unless_drained_with_close, a_refused_spec_changes_nothing, an_established_connection_survives_an_apply_that_keeps_its_inbound, shutdown_ends_pending_grace_timers. |
| Tracking | supervisor/tests/tracking.rs |
The metadata the connector files: outbound tag and rule() per mux sub-flow and per UDP sub-link. |
| App end to end | app/tests/integration/e2e_route_context.rs |
inbound_tag_selects_the_route, source_cidr_matches_the_client_address, network_separates_tcp_from_udp, each against a config where the matcher is the only thing between the traffic and blackhole. |
| App end to end | app/tests/integration/e2e_sniff.rs |
A VLESS inbound behind an Xray client: an HTTP Host and a TLS SNI route an IP-addressed flow, a non-matching name and an unsniffable payload fall to the default, and sniffing = false stops the match. A blocked flow never answers, so a 10 s timeout on the round trip is the test’s “dropped” signal rather than a failure. |