Tracking flows, stats and speed limits
Source files: 47 · checked against Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/metered.rsEtemenanki/supervisor/src/track/pace.rsEtemenanki/supervisor/src/track/sampler.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/topology/outbound/dns.rsEtemenanki/supervisor/src/topology/balancer.rsEtemenanki/supervisor/src/topology/plane.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/user.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/usage.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/build/inbound.rsEtemenanki/supervisor/src/topology/inbound/mod.rsEtemenanki/supervisor/src/lib.rsEtemenanki/supervisor/Cargo.tomlEtemenanki/concepts/src/runtime.rsEtemenanki/concepts/src/relay.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/mux/mod.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/protocols/src/vless/core.rsEtemenanki/protocols/src/vmess/core.rsEtemenanki/protocols/src/trojan/core.rsEtemenanki/protocols/src/ss_legacy/core.rsEtemenanki/protocols/src/ss_2022/core.rsEtemenanki/protocols/src/http/core.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/socks/server.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/app/src/lower.rsEtemenanki/app/src/instance.rsEtemenanki/webclient/src/router.rsEtemenanki/webclient/src/tracking.rsEtemenanki/ffi/src/proxy.rsEtemenanki/supervisor/tests/unit/track.rsEtemenanki/supervisor/tests/tracking.rsEtemenanki/supervisor/tests/hot_swap.rsEtemenanki/supervisor/benches/tracking.rsEtemenanki/webclient/tests/tracking.rskatana/src/lower/mod.rskatana/src/manager/node.rs
The supervisor tracks every flow it routes. A flow is one routed outbound: a stream, or one sub-link of a UDP association toward one outbound version. A plain proxied TCP connection is one flow; a mux connection or a Hysteria 2 connection opens many, and so does a gRPC connection, each of whose HTTP/2 streams is a session of its own. The connector wraps whatever it opens in a Metered stream or a MeteredDatagram link. The wrapper registers the flow in a concurrent map, counts its payload in both directions, fails it once it is killed, and paces it while its user’s token bucket is in debt. A sampler task reads every flow’s counters once per tick and publishes totals and rates per inbound, per outbound and per user.
This page is for contributors who change supervisor/src/track/, the connector or the UDP fan-out, and for front-end authors who read the tracker. It covers the registry, the wrappers, the kill switch, the lifecycle events, the sampler and StatsSnapshot, per-user speed limits, and the benchmark behind the registry’s design. The wire bytes each session counts, and the usage ledger the sampler reconciles, are on Per-user usage accounting. How a flow is routed, and the drain guard inside each wrapper, are on The plane: routing each flow.
Responsibilities
Section titled “Responsibilities”The tracking code (supervisor/src/track/):
- registers every flow the connector opens, with metadata fixed at open (
FlowMeta), and deregisters it when its wrapper drops; - counts each flow’s payload on counters that belong to that flow alone;
- closes one flow, or every flow a predicate selects, from outside the task that relays it;
- announces every open and close on a bounded broadcast channel;
- samples every flow once per tick into a
StatsSnapshoton awatchchannel, and reconciles the usage ledger in the same tick; - paces the flows of a speed-limited user with that user’s one token bucket, shared by all their flows in both directions.
It does not:
- count wire bytes. A session’s
Wiredoes that, and the ledger bills it (Per-user usage accounting). Flow counters are payload. The one exception is a charged datagram flow (a SOCKSUDP ASSOCIATEsub-link), whose wrapper also adds its payload to its session’sWire; seechargeunder Where flows are metered. - close sessions.
Supervisor::closecloses whole sessions through the actor, by session, user, inbound or all (Users, principals and sessions). A kill ends one flow, and the session and its other flows run on. - route or drain. A flow is routed before it is wrapped, and the
Guardedlink inside the wrapper fails the flow withthe outbound this flow was opened on was drainedonce a drain policy closes its outbound version (The plane). - see what the supervisor dials for itself.
RoutedDialer(supervisor/src/topology/outbound/dns.rs) routes the DNS resolver’s connections to its servers and opens them withTarget::connect_streamdirectly, and a balancer’s health probe (supervisor/src/topology/balancer.rs→probe) is a bare TCP connect. Neither passes the connector, so neither is registered or counted here.
The module’s doc states its design rule as the data plane never waits on an observer, and defines it by two properties:
- Each flow’s counters are its own, written only by the task that relays the flow, so they are uncontended. Nothing per tag or per user is counted per byte; the sampler aggregates the flows’ counters once per tick.
- Lifecycle events go out on a bounded broadcast channel. A slow subscriber lags and loses events instead of holding up a flow’s open or close.
The rule is about counting and events, not about every shared structure. Opens and closes share the registry with a running scan, which slows them (see The registry and why scc), and open takes the pacers’ read lock for a user’s flow, which set_speed_limits holds for writing during its pass over every pacer.
Who reads it
Section titled “Who reads it”Supervisor::tracker() hands out the Tracker without going through the actor; the crate docs list it, with usage, as what “is observable without going through the actor”. Tracker‘s own doc says “Nothing here depends on the config”. The actor builds it once, in Actor::new, and no apply replaces it; the only thing in it that follows the spec is the users’ speed limits.
| Reader | Uses | Where |
|---|---|---|
| etemenanki-app’s REST API | flows, sessions, kill, kill_where, subscribe, stats |
webclient/src/tracking.rs: GET /v1/connections, DELETE /v1/connections/{flow}, DELETE /v1/connections?outbound=…&inbound=…&session=…&user=…, GET /v1/connections/events, GET /v1/traffic |
| The FFI library | stats, flows, kill |
ffi/src/proxy.rs → traffic, connections (sorted by id, oldest first), close_connection |
| katana | subscribe |
src/manager/node.rs → spawn_audit, which records audit hits from FlowEvent::Opened |
A few REST details follow the tracker’s semantics directly. GET /v1/connections sorts the flows by id. DELETE /v1/connections/{flow} is Tracker::kill: it answers 204 whenever a flow with that id is registered (killed before or not), or 404 with no live flow <id> when none is, for example once the flow’s wrapper has dropped. The selector of DELETE /v1/connections is kill_where over the fields given (all flows when none is); its user field matches the user’s label and only flows that have a user, and it rejects an unknown query key with 400 (deny_unknown_fields). The answer is {"closed": n}, the count kill_where returned.
The REST API page describes the endpoints, the FFI page describes the calls, and the katana audit page describes the audit rules.
Key types
Section titled “Key types”Tracker
Section titled “Tracker”const EVENT_CAPACITY: usize = 1024;
#[derive(Clone)]pub struct Tracker { inner: Arc<Inner>,}
struct Inner { flows: scc::HashMap<FlowId, Arc<FlowEntry>>, next: AtomicU64, events: broadcast::Sender<FlowEvent>, /// Flows that closed, for the sampler to fold their last bytes in. closed: mpsc::UnboundedSender<Arc<FlowEntry>>, stats: watch::Receiver<Arc<StatsSnapshot>>, sessions: Arc<Sessions>, /// Each user's speed limit, shared by their flows. pacers: parking_lot::RwLock<HashMap<UserKey, Arc<Pacer>>>,}
impl Tracker { pub(crate) fn new(sessions: Arc<Sessions>) -> (Self, sampler::Sampler); pub fn flows(&self) -> Vec<Arc<FlowEntry>>; pub(crate) fn for_each_flow(&self, visit: impl FnMut(&Arc<FlowEntry>)); pub fn flow(&self, id: FlowId) -> Option<Arc<FlowEntry>>; pub fn sessions(&self) -> Vec<SessionStats>; pub fn kill(&self, id: FlowId) -> bool; pub fn kill_where(&self, select: impl FnMut(&FlowEntry) -> bool) -> usize; pub fn subscribe(&self) -> broadcast::Receiver<FlowEvent>; pub(crate) fn set_speed_limits(&self, limits: &HashMap<UserKey, NonZeroU64>); pub fn stats(&self) -> watch::Receiver<Arc<StatsSnapshot>>; pub(crate) fn meter_stream<S>(&self, meta: FlowMeta, stream: S) -> Metered<S>; pub(crate) fn meter_datagram<D>( &self, meta: FlowMeta, link: D, charge: Option<Arc<Wire>>, ) -> MeteredDatagram<D>;
fn pacer(&self, key: UserKey) -> Arc<Pacer>; fn open(&self, meta: FlowMeta, charge: Option<Arc<Wire>>) -> FlowHandle;}Tracker is a cheap handle: cloning it clones one Arc. The module re-exports Metered, MeteredDatagram, StatsSnapshot, TagStats and UserStats, so a front end names everything as etemenanki_supervisor::track::….
Field of Inner |
Type | Holds |
|---|---|---|
flows |
scc::HashMap<FlowId, Arc<FlowEntry>> |
The live-flow registry. Why scc: see The registry and why scc. |
next |
AtomicU64, starting at 1 |
The next flow id. open takes it with fetch_add(1, Relaxed), so ids are unique for the supervisor’s lifetime and increase in the order open takes them. Taking the id and inserting the entry are separate steps, so two concurrent opens may reach the registry in the opposite order to their ids. |
events |
broadcast::Sender<FlowEvent> of capacity EVENT_CAPACITY = 1024 |
Open and close events. The tracker keeps only the sender; every receiver comes from subscribe. |
closed |
mpsc::UnboundedSender<Arc<FlowEntry>> |
The closed queue: each flow that deregisters is sent here so the sampler can count what it moved since the last tick. The sampler owns the receiver (graveyard) and empties it every tick. |
stats |
watch::Receiver<Arc<StatsSnapshot>> |
The last published snapshot. The sampler owns the sender. |
sessions |
Arc<Sessions> |
The session registry. sessions() reads it, and the sampler reaches the usage ledger through it. |
pacers |
parking_lot::RwLock<HashMap<UserKey, Arc<Pacer>>> |
One token bucket per user key. An unlimited pacer costs one atomic load in poll_ready and one in charge. |
| Method | What it does | Cost |
|---|---|---|
flows() |
Every live entry, cloned into a Vec sized by the map’s len(), in no particular order |
One pass over the map (iter_sync) |
flow(id) |
The live entry id, if it is still registered |
One read_sync |
sessions() |
Sessions::stats(): every live session with its wire bytes, sorted by session id (Users, principals and sessions) |
One pass over the sessions under the session registry’s lock, then one Wire read per session after releasing it, because reading a QUIC session’s bytes takes its connection’s lock and opens and admissions must not wait on that |
kill(id) |
Kills flow id. Returns true whenever id is in the registry, whether or not it was killed before, and false once its wrapper has dropped. |
One read_sync |
kill_where(select) |
Kills every live flow that select picks and that is not already killed, and returns how many it killed |
One pass; select runs inside iter_sync |
subscribe() |
A receiver of every event sent from now on | None |
stats() |
A clone of the watch receiver |
None |
set_speed_limits(limits) |
Sets every pacer’s rate; see Speed limits | One pass over pacers under the write lock |
meter_stream, meter_datagram |
Register a flow and wrap its link | One insert_sync and one event |
for_each_flow(visit) |
Visits every live flow without collecting; the sampler’s scan | One pass |
FlowMeta
Section titled “FlowMeta”#[derive(Debug, Clone)]pub(crate) struct FlowMeta { pub session: Option<SessionId>, pub inbound: CompactString, pub principal: Arc<Principal>, pub source: Option<IpAddr>, pub destination: Destination, pub sniffed: Option<CompactString>, pub outbound: OutboundId, pub rule: Option<RuleId>, pub plane_epoch: u64,}
impl FlowMeta { pub(crate) fn of( flow: &Flow, session: Option<SessionId>, inbound: &CompactString, listener_source: Option<IpAddr>, outbound: OutboundId, rule: Option<RuleId>, plane_epoch: u64, ) -> Self;}FlowMeta is what a flow is, fixed when it opens. of builds it from the routed Flow and the connector’s context:
| Field | Taken from |
|---|---|
session |
The connector’s session id. None when the connector has no session: a TUN device’s flows, since a device admits no users and has no sessions, and a Hysteria 2 connection that finished its handshake as its listener was stopped, which run_hysteria_inbound (supervisor/src/serve.rs) gives a session-less connector and a token that is already cancelled |
inbound |
The inbound tag in the connector’s FlowContext |
principal |
flow.user.user_data: the Principal the inbound admitted, with its user key and label, or Principal::anonymous() |
source |
flow.source, or the listener’s peer address when the flow names none |
destination |
flow.destination; for a UDP sub-link, the destination of the packet that opened it |
sniffed |
The domain sniffing recovered, if any. Always None for a UDP sub-link, whose Flow comes from Flow::toward (protocols/src/flow.rs), which sets sniffed: None |
outbound |
The OutboundId of the version the flow is opened on. A balancer resolves to its member first, so the flow names the member, not the balancer. |
rule |
The index of the rule that matched in that plane’s RouteSpec::rules; None for the default route and for a DNS query the supervisor’s own DNS service answers |
plane_epoch |
The epoch of the plane that routed it |
FlowEntry
Section titled “FlowEntry”pub struct FlowEntry { id: FlowId, meta: FlowMeta, started: Instant, // Written by the task relaying the flow alone: uncontended. up: AtomicU64, down: AtomicU64, // What the sampler has counted so far; only the sampler touches these. sampled_up: AtomicU64, sampled_down: AtomicU64, kill: AtomicBool, read_waker: AtomicWaker, write_waker: AtomicWaker,}
impl FlowEntry { pub fn id(&self) -> FlowId; pub fn session(&self) -> Option<SessionId>; pub fn inbound(&self) -> &str; pub fn user(&self) -> Option<UserKey>; pub fn user_label(&self) -> &str; pub fn source(&self) -> Option<IpAddr>; pub fn destination(&self) -> &Destination; pub fn sniffed(&self) -> Option<&str>; pub fn outbound(&self) -> &OutboundId; pub fn rule(&self) -> Option<RuleId>; pub fn plane_epoch(&self) -> u64; pub fn started(&self) -> Instant; pub fn up(&self) -> u64; pub fn down(&self) -> u64; pub fn kill(&self); pub fn killed(&self) -> bool;
fn take_sample(&self) -> (u64, u64);}One FlowEntry per flow, behind an Arc shared by the registry, the wrapper, any event still queued, and any caller that holds it. The accessors read meta: user() is the principal’s user key (None for nobody), user_label() its label (empty for nobody). started is a std::time::Instant taken when the flow registered, which is after its dial completed.
| Field | Written by | Ordering | Meaning |
|---|---|---|---|
up |
The relaying task, after each successful write or send | Relaxed |
Payload bytes toward the destination so far |
down |
The relaying task, after each successful read or receive | Relaxed |
Payload bytes back from the destination so far |
sampled_up, sampled_down |
The sampler only, in take_sample |
Relaxed |
The watermark: how much of up and down the sampler has already counted |
kill |
kill() |
Release store, Acquire load |
Set once; never cleared |
read_waker, write_waker |
The wrapper, whenever a poll returns Pending |
futures::task::AtomicWaker |
The tasks to wake on a kill |
kill() stores true and then wakes both wakers. take_sample() swaps each watermark for the current counter and returns the differences; it is private, and only the sampler calls it.
FlowEntry has a hand-written Debug. It prints id, meta (the derived Debug of FlowMeta), up, down and killed, reading the counters and the flag through their accessors; started, the watermarks and the wakers are left out.
FlowEvent
Section titled “FlowEvent”#[derive(Clone)]pub enum FlowEvent { Opened(Arc<FlowEntry>), Closed(Arc<FlowEntry>),}Both variants carry the entry itself, not a copy of its counters. An Opened entry keeps counting after the event was sent; a Closed entry carries its final counts, since its wrapper is gone. The hand-written Debug prints only the variant and the flow id.
FlowHandle
Section titled “FlowHandle”struct FlowHandle { entry: Arc<FlowEntry>, tracker: Arc<Inner>, /// The session wire the flow's payload also counts toward, if any. charge: Option<Arc<Wire>>, /// The user's speed limit, charged with the payload; `None` for a flow /// no user was admitted for. pacer: Option<Arc<Pacer>>,}
impl FlowHandle { fn check(&self) -> Option<io::Error>; fn add_up(&self, n: usize); fn add_down(&self, n: usize);}
impl Drop for FlowHandle { /* deregister */ }
fn killed() -> io::Error { io::Error::new(io::ErrorKind::ConnectionAborted, "the flow was killed")}The private registration token each wrapper owns. check returns the kill error once the flow is killed. add_up and add_down add n to the entry’s counter, to the charged Wire when there is one (wire.add(n, 0) or wire.add(0, n)), and to the pacer when there is one. Dropping the handle is what deregisters the flow.
Metered and MeteredDatagram
Section titled “Metered and MeteredDatagram”pub struct Metered<S> { inner: S, flow: FlowHandle, read_pause: Option<Pin<Box<Sleep>>>, write_pause: Option<Pin<Box<Sleep>>>,}
impl<S> Metered<S> { pub fn entry(&self) -> &FlowEntry;}impl<S: AsyncRead + Unpin> AsyncRead for Metered<S> { /* poll_read */ }impl<S: AsyncWrite + Unpin> AsyncWrite for Metered<S> { /* poll_write, poll_flush, poll_shutdown */}
pub struct MeteredDatagram<D> { inner: D, flow: FlowHandle, read_pause: Option<Pin<Box<Sleep>>>, write_pause: Option<Pin<Box<Sleep>>>,}
impl<D> MeteredDatagram<D> { pub fn entry(&self) -> &FlowEntry;}impl<D: DatagramLink<Addr = Destination>> DatagramLink for MeteredDatagram<D> { type Addr = Destination; /* poll_send_to, poll_recv_from */}Both constructors are pub(super): only Tracker::meter_stream and Tracker::meter_datagram build them. read_pause and write_pause are the pacing timers, one per direction, allocated only while the user’s bucket is in debt.
const MAX_WAIT: Duration = Duration::from_secs(1);const MIN_WAIT: Duration = Duration::from_millis(1);
pub(crate) struct Pacer { /// The current rate; 0 is unlimited, the fast path that takes no lock. limited: AtomicU64, state: Mutex<PaceState>,}
struct PaceState { /// 0 while unlimited. rate: u64, /// Negative while in debt. tokens: f64, last: Instant,}
impl Pacer { pub(crate) fn unlimited() -> Self; pub(crate) fn set_rate(&self, rate: Option<NonZeroU64>); pub(crate) fn charge(&self, n: usize); pub(crate) fn poll_ready( &self, cx: &mut Context<'_>, pause: &mut Option<Pin<Box<Sleep>>>, ) -> Poll<()>;}One user’s bucket: rate bytes per second with a burst of one second of rate. Mutex is parking_lot::Mutex, and Instant and Sleep are Tokio’s, so the bucket runs on Tokio’s clock and the tests can pause it. Speed limits walks through each method.
Sampler and StatsSnapshot
Section titled “Sampler and StatsSnapshot”pub const DEFAULT_SAMPLE_INTERVAL: Duration = Duration::from_secs(1);
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]pub struct TagStats { pub up: u64, pub down: u64, pub up_rate: u64, pub down_rate: u64, pub flows: usize,}
#[derive(Debug, Clone, PartialEq, Eq)]pub struct UserStats { pub user: UserKey, pub label: CompactString, pub stats: TagStats,}
#[derive(Debug, Clone, Default, PartialEq, Eq)]pub struct StatsSnapshot { pub tick: u64, pub interval: Duration, pub total: TagStats, pub inbounds: BTreeMap<CompactString, TagStats>, pub outbounds: BTreeMap<CompactString, TagStats>, pub users: Vec<UserStats>,}
pub(crate) struct Sampler { tracker: Tracker, graveyard: mpsc::UnboundedReceiver<Arc<FlowEntry>>, publish: watch::Sender<Arc<StatsSnapshot>>, totals: Totals,}
impl Sampler { pub(crate) async fn run(self, interval: Duration, stop: CancellationToken); pub(crate) fn tick(&mut self, interval: Duration) -> StatsSnapshot;}| Field | Meaning |
|---|---|
TagStats::up, down |
Payload toward destinations and back, in total since the supervisor started, closed flows included |
TagStats::up_rate, down_rate |
Bytes per second over the last tick: the bytes counted in this tick divided by interval |
TagStats::flows |
Flows live at the tick |
StatsSnapshot::tick |
Ticks so far; 0 in the default snapshot the channel holds before the first tick |
StatsSnapshot::interval |
The time the rates are over: the time actually measured since the previous tick, not the configured interval |
StatsSnapshot::total |
Every flow together, anonymous ones included |
StatsSnapshot::inbounds |
By inbound tag, every inbound that has carried a flow since the supervisor started |
StatsSnapshot::outbounds |
By outbound tag, every outbound that has carried a flow. The supervisor’s own DNS service is the empty tag "" (its internal OutboundId has an empty tag, which no spec tag can be). The blackhole target of Plane::empty(), the plane the supervisor holds before its first apply, also has the empty tag, at version 0. |
StatsSnapshot::users |
The users with a live flow or traffic in the last tick, sorted by UserKey. Anonymous flows count toward total, inbounds and outbounds, never here. |
UserStats::label |
The label the protocols know the user by, taken from the first of their flows the sampler counted and never updated |
DEFAULT_SAMPLE_INTERVAL lives in a crate-private module; front ends see it as the default of SupervisorBuilder::sample_interval, documented as “one second when not set”. etemenanki-app (app/src/instance.rs), the FFI library (ffi/src/proxy.rs) and katana (src/manager/node.rs) all keep the default.
Identifiers
Section titled “Identifiers”The ids come from supervisor/src/entity/id.rs. FlowId, SessionId, UserKey and RuleId are Copy newtypes over an integer, with new and get (both const fn); OutboundId is a struct of a tag and a version:
| Type | Wraps | Display |
Meaning here |
|---|---|---|---|
FlowId |
u64 |
flow#N |
One flow, unique within a supervisor’s lifetime |
SessionId |
u64 |
session#N |
The session a flow belongs to |
UserKey |
u64 |
none | The supervisor’s index for a front end’s user id; keys the pacers and UserStats |
OutboundId |
tag: CompactString, version: u64 |
tag@vN |
The outbound version a flow is opened on |
RuleId |
u32 |
none | A rule’s index in the routing plane’s RouteSpec::rules |
Where flows are metered
Section titled “Where flows are metered”On the data plane, exactly two places register flows, and both are on the path every inbound dials through. The only other caller of the tracker’s registration is the benchmark helper track::bench::metered (The benchmark).
| Where | Opens | Wrapper | Registered | charge |
|---|---|---|---|---|
supervisor/src/connector.rs → AppConnector::connect, for a TCP flow |
One stream per connect: a connection’s single flow, each mux sub-flow, each Hysteria 2 proxy stream, each TUN TCP flow, each SOCKS CONNECT |
Metered<Guarded<OutboundStream>> |
When Target::connect_stream resolves Ok |
None: the session’s wire is counted on the client side |
supervisor/src/topology/outbound/udp_fanout.rs → FanOutLink::poll_opening, for a UDP association |
One sub-link per outbound version the association’s packets are routed to | MeteredDatagram<Guarded<OutboundDatagram>> |
When the sub-link’s connect_datagram resolves Ok |
The session’s Wire when the connector was built with charging_datagrams() (SOCKS UDP ASSOCIATE), otherwise None |
A UDP flow is not dialled at all: connect returns a FanOutLink at once, and the fan-out routes packet by packet (Outbounds, UDP fan-out and balancers). Each sub-link it opens is a flow of its own, registered under the FlowScope the connector gave it (the tracker, the session id and the optional charge). A sub-link that the fan-out drops, because it failed, because it was killed, or because it was the least recently used one past MAX_SUBS = 64, drops its wrapper and so deregisters its flow. The table keeps the most recently sent-to sub-link at the back, so the eviction removes index 0 and moves the receive cursor next back by one, stopping at 0. A sub-link whose connect_datagram fails is never registered; poll_opening logs udp fan-out: opening <tag>@v<version> failed: <error> at debug.
The wrapper sits outside Guarded, so every read, write, send or receive checks the kill flag first, then the pacer, then the drain guard, then the link itself.
flowchart LR conn["AppConnector::connect"] tcp["Target::connect_stream"] fan["FanOutLink"] sub["poll_opening: connect_datagram"] ms["meter_stream"] md["meter_datagram"] open["Tracker::open"] map["scc registry"] ev["broadcast: FlowEvent"] wrap["Metered or MeteredDatagram"] conn -->|"TCP"| tcp conn -->|"UDP"| fan fan --> sub tcp -->|"Ok"| ms sub -->|"Ok"| md ms --> open md --> open open -->|"insert_sync"| map open -->|"Opened"| ev open -->|"FlowHandle"| wrap
Registering: Tracker::open
Section titled “Registering: Tracker::open”- When the principal has a user key, look up the user’s pacer (
pacer(key)): a hit under the read lock, or else a new unlimited pacer inserted under the write lock. Every user’s flow gets a pacer even without a limit, so a limit set later reaches flows that are already open. - Take the next id from
next. - Build the entry: counters and watermarks at 0,
killfalse,startednow. insert_syncit into the registry. The result is ignored; ids never repeat.- Send
FlowEvent::Opened. A send error, which only means nobody is subscribed, is ignored. - Return the
FlowHandle, which the wrapper keeps.
The FlowMeta of a TCP flow is built before the dial starts, but the flow registers only when the dial succeeds. connect loads the plane once per call (“Loaded once for this flow: a route change reaches the next flow”), so the route, the balancer’s pick and the plane epoch in FlowMeta are the ones fixed at that load. A dial in progress or a failed dial is never in the registry, sends no event and counts nothing.
Deregistering: FlowHandle::drop
Section titled “Deregistering: FlowHandle::drop”remove_syncthe entry from the registry.- Send the entry on the closed queue, so the sampler counts what it moved since the last tick. The send fails silently once the sampler has stopped.
- Send
FlowEvent::Closedwith the entry, now at its final counts.
The wrapper drops when the runtime drops the outbound (after an error, a close effect or the end of the connection), when the fan-out drops a sub-link, or when the connection’s task is dropped. There is no other way out of the registry: an entry lives exactly as long as its wrapper.
What is counted
Section titled “What is counted”The counters hold payload between the protocol core and the outbound, not wire bytes:
| Wrapper method | Counts | After |
|---|---|---|
Metered::poll_write |
add_up(n), n being what the inner write accepted |
Ready(Ok(n)) only |
Metered::poll_read |
add_down of how much buf.filled() grew |
Ready(Ok(())) only; an end of stream adds 0 |
MeteredDatagram::poll_send_to |
add_up(n) |
Ready(Ok(n)) only |
MeteredDatagram::poll_recv_from |
add_down of how much buf.filled() grew |
Ready(Ok(_)) only |
poll_flush, poll_shutdown |
Nothing |
For a proxy outbound, that is the payload before the client-side protocol frames and encrypts it. So a flow’s counters, a session’s wire bytes and the upstream’s own bytes all differ; Per-user usage accounting compares the first two.
Killing a flow
Section titled “Killing a flow”FlowEntry::kill (or Tracker::kill, Tracker::kill_where) closes one flow from outside the task that relays it. It is two steps: store true in kill with Release ordering, then wake read_waker and write_waker. From then on, every poll of the wrapper in either direction returns Err with io::ErrorKind::ConnectionAborted and the text the flow was killed, before it touches the pacer or the link. This includes poll_flush and poll_shutdown, so a runtime that is waiting only on a shutdown also learns of the kill.
The race with a pending poll
Section titled “The race with a pending poll”Every wrapper poll goes through one of two helpers in supervisor/src/track/metered.rs:
guarded(flow, waker, cx, op): return the kill error if the flow is killed; otherwise runop. IfopreturnsPending, registercx’s waker inwaker, then check the flag again, and return the kill error if it is set now.paced(flow, waker, pause, cx, op): return the kill error if the flow is killed. If the flow has a pacer andPacer::poll_readyisPending, register the waker and check the flag again, exactly as above. Otherwise continue withguarded.
Registering before the second check closes the race: a kill that lands between the first check and the registration is either seen by the second check, or finds the waker registered and wakes the task. poll_read and poll_recv_from park on read_waker; poll_write, poll_send_to, poll_flush and poll_shutdown park on write_waker. A flow parked on its user’s debt is woken by the kill too, not only when its pacing timer fires (killing_a_paced_flow_fails_its_pending_write).
What the connection does with the error
Section titled “What the connection does with the error”A kill takes effect at the wrapper’s next poll, and what follows depends on who holds the wrapper:
sequenceDiagram participant C as caller participant E as FlowEntry participant R as runtime task participant M as Metered participant K as protocol core C->>E: kill E->>E: kill flag set with Release E->>R: wake read_waker and write_waker R->>M: poll_read or poll_write M->>R: Err ConnectionAborted, the flow was killed R->>R: drop the key and its link M->>E: FlowHandle dropped, flow deregistered R->>K: Event::OutboundError for the key K->>R: Close the sub-flow, or ShutdownTransport and Finish
The runtime treats the error like any failed read or write: it drops the key’s slot, which drops the wrapper, and hands the core Event::OutboundError (The server runtime). The core decides the rest:
| Flow | Receives the error | Result | Debug line |
|---|---|---|---|
| The one TCP flow of an HTTP, Shadowsocks, Shadowsocks 2022, Trojan, VLESS or VMess connection | The connection’s runtime, then the core | on_outbound_gone: the core pushes ShutdownTransport and Finish (The server runtime) |
http: outbound failed: the flow was killed; shadowsocks: outbound gone: …, shadowsocks-2022: outbound gone: …, trojan: outbound gone: …, vless: outbound gone: …, vmess: outbound gone: … |
One sub-flow of a VLESS, VMess or Trojan mux connection (a Trojan client reaches the demultiplexer with a CONNECT to the mux address, is_mux_destination in protocols/src/mux/mod.rs) |
The carrier’s runtime, then the core | Demux::on_outbound_gone: the demultiplexer sends an End frame for that sub-flow and pushes Close for its key; the carrier and the other sub-flows run on, and new sub-flows still open |
vless: outbound gone: the flow was killed, vmess: … or trojan: … |
| One Hysteria 2 proxy stream | That stream’s runtime, then Hy2StreamCore |
on_outbound_gone on that stream’s relay: ShutdownTransport and Finish for that stream alone; the QUIC connection (the session) and its other streams run on |
hysteria2: outbound failed: the flow was killed |
| One TUN TCP flow | That flow’s runtime, then PassthroughCore |
on_outbound_gone: ShutdownTransport and Finish for that TCP flow |
None |
A SOCKS CONNECT |
The SOCKS driver’s relay | The relay returns the error and the connection closes | socks connection from Some(<client IP>) ended: the flow was killed, from serve_connection |
| One UDP sub-link, from any inbound | FanOutLink |
The fan-out drops that sub-link; the association and its other sub-links run on, and a later packet routed there opens a new flow with a new id. The core receives no event. | udp fan-out: sending on <tag>@v<version> failed: the flow was killed, or udp fan-out: sub-link <tag>@v<version> ended: the flow was killed |
All of these are debug lines. The fan-out drops the packet whose send failed and reports it as sent, as a lossy link would.
A kill is not a session close. To end whole connections, a front end calls Supervisor::close with a Selector, which goes through the actor and cancels session tokens (Users, principals and sessions).
kill against kill_where
Section titled “kill against kill_where”Tracker::kill(id) runs kill() under read_sync and returns whether the id was registered, so it returns true again for a flow that was already killed but whose wrapper has not dropped yet. kill_where skips flows that are already killed and counts only those it kills itself, so calling it twice with the same predicate returns 0 the second time (a_selective_kill_leaves_the_other_flows).
Lifecycle events
Section titled “Lifecycle events”Every open and close is sent on one tokio::sync::broadcast channel of capacity EVENT_CAPACITY = 1024.
Openedis sent after the entry is in the registry, so a subscriber can look it up withflow(id)until it closes.Closedis sent after the entry left the registry and was queued for the sampler, with its final counts.- For one flow,
Openedis always sent beforeClosed: the first is sent byopen, the second by the handle’sDrop.
A receiver from subscribe() sees only events sent after it subscribed, so it can receive a Closed for a flow whose Opened it never saw. Sending never waits: the channel keeps the last 1024 events, and a receiver that falls further behind gets RecvError::Lagged(n) on its next recv, with n the number of events it missed, and then continues from the oldest event still kept. With no receiver at all, the send fails and the event is dropped at once; the error is ignored. The REST API turns a lag into a lagged event carrying missed.
The sampler
Section titled “The sampler”The sampler turns the flows’ counters into a StatsSnapshot once per tick. It is the only reader that marks bytes as counted, through each flow’s watermark, so every byte a flow moves is counted in exactly one tick.
Starting and stopping
Section titled “Starting and stopping”Tracker::new creates the sampler together with the tracker, and Actor::new keeps it in Actor::sampler until SupervisorBuilder::start runs it. SupervisorBuilder’s Default sets sample_interval to DEFAULT_SAMPLE_INTERVAL, and Supervisor::start(spec) is Supervisor::builder().start(spec), so a supervisor started without a builder samples once a second.
-
startasserts that the interval is not zero, before anything else. A zero interval panics withthe sample interval must not be zero;SupervisorBuilder::sample_intervaldocuments the panic. -
startapplies the first spec, then spawnssampler.run(sample_interval, background_stop)withtokio::spawnand keeps itsJoinHandleinActor::background. -
At shutdown,
Actor::shutdownworks in this order:- stops every listener;
- cancels every balancer’s probe token;
- closes the task tracker that every accept loop and connection runs on, and waits for it up to the grace;
- cancels the root token, then waits for the task tracker again;
- clears the listeners;
- cancels
background_stopand awaits every background task, the sampler among them.
The sampler’s last tick therefore comes after every accept loop and connection on the task tracker has ended.
Supervisor::shutdown(grace) runs this through the actor. When every Supervisor handle is dropped without a shutdown, the actor’s command loop ends and calls shutdown(Duration::ZERO) itself, which stops the sampler the same way. check() also builds an actor and so a sampler, which it never runs.
The loop: Sampler::run
Section titled “The loop: Sampler::run”let mut ticker = tokio::time::interval_at(Instant::now() + interval, interval);ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);- The first tick comes one interval after the start, not at once (
the_sampler_publishes_once_per_interval_and_once_when_stoppedsees tick 1 at exactly 1 s and tick 2 at 2 s). MissedTickBehavior::Delay: a tick delayed by a busy runtime is not made up in a burst. The next tick’s rates cover the longer interval instead, becauseintervalis measured (now - last, on Tokio’s clock) rather than assumed.- Each iteration waits on
stop.cancelled()andticker.tick()in aselect!, then ticks either way. When the stop won, it publishes that last tick and returns. The last snapshot’s rates cover the time since the previous tick. - The tick itself runs in
tokio::task::spawn_blocking: the sampler is moved into the closure and handed back with the snapshot. A tick scans every live flow and reconciles every active user’s live sessions, and reading a Hysteria 2 session’s bytes takes its QUIC connection’s lock, so none of it runs on an async worker. A panic inside the tick is re-raised in the sampler’s task withresume_unwind. - The snapshot is published with
watch::Sender::send_replace(Arc::new(snapshot)), which succeeds whether or not anyone is watching.
One tick: Sampler::tick
Section titled “One tick: Sampler::tick”flowchart TB rec["Ledger::reconcile"] reset["tick += 1, reset per-tick counts"] scan["for_each_flow: take_sample, count as live"] drain["graveyard.try_recv until empty: take_sample, count as closed"] snap["snapshot: rates over the measured interval"] pub["publish on the watch channel"] rec --> reset --> scan --> drain --> snap --> pub
- Reconcile the usage ledger.
self.tracker.inner.sessions.ledger().reconcile()moves each live session’s new wire bytes into its user’s pending total, so a usage take only drains totals (Per-user usage accounting). - Start the tick.
tickgoes up by one, and every group’s per-tick bytes and live-flow count go back to zero. Totals are kept. - Scan the live flows. For each registered flow,
take_samplereturns the bytes since its watermark and moves the watermark up;countadds them to the total, to the flow’s inbound group, to its outbound group (by tag, not version) and, when the flow has a user, to that user’s group, and counts the flow as live in each. - Drain the closed queue. For each flow that deregistered since the last drain,
take_samplereturns what it moved since its last sample, andcountadds it without counting the flow as live. - Snapshot. Each group becomes a
TagStats. Its rates come fromper_second(bytes, interval), which computesbytes * 1_000_000_000 / interval_nanosinu128, givesu64::MAXwhen the result does not fit au64(unwrap_or(u64::MAX)), and 0 for a zero interval. Users are filtered to those with a live flow or bytes in this tick, and sorted by key.
The order of steps 3 and 4, and the watermarks, are what make the count exact. A flow that closes during the scan was either visited, and then its queue entry gives only what it moved since, or it was not, and then all its uncounted bytes come from its queue entry. A flow whose entry reaches the queue after the drain is counted by the next tick. Either way the watermark swap hands each byte to exactly one tick (the_sampler_counts_every_byte_once_by_inbound_outbound_and_user).
Aggregates
Section titled “Aggregates”Totals holds the tick counter and one Group for everything, one per inbound tag, one per outbound tag and one per user key (with the label):
struct Group { up: u64, down: u64, tick_up: u64, tick_down: u64, flows: usize,}Groups are created on first use, so inbounds and outbounds also list tags that have carried a flow and are idle now, at zero rates. A tag is looked up by &str first, so a known tag costs no allocation. A user’s label is stored when the user’s group is created, from the first of their flows the sampler counts, and is never updated.
StatsSnapshot counts payload per flow; usage counts wire bytes per session. The two cannot be derived from each other; the table on Per-user usage accounting sets them side by side.
Speed limits
Section titled “Speed limits”A user’s speed limit is UserSpec::speed_limit: Option<NonZeroU64>, in payload bytes per second, shared by every flow of that user in both directions, TCP and UDP. None is unlimited. Changing it does not change how the user is admitted, so the user’s sessions stay open (a_speed_limit_change_keeps_the_users_sessions).
Publishing the limits
Section titled “Publishing the limits”The limits reach the tracker only when an apply or a user edit commits:
Actor::commit(every apply, including the first one instart) andActor::edit_users(afterset_users,upsert_userorremove_user) callpublish_speed_limits(&spec)right after adopting the new user keys.publish_speed_limitswalks every user of every user set in the new spec. It skips a user with nospeed_limitor no user key, and a user in several sets gets the smallest of their limits.Tracker::set_speed_limits(&limits)takes the pacers’ write lock, sets every existing pacer to its user’s limit or to unlimited when the user is not inlimits, then creates a pacer, set to its rate, for each limited user that has none yet.
Only at commit, because a refused apply or a refused user edit must not change the limits that are live. Every live flow already holds its user’s Arc<Pacer>, so a new rate reaches it at its next poll, without reopening anything. The crate docs put it as: a user’s limit “takes effect on live flows at their next read or write, without ending their sessions”.
Each front end sets the limits. etemenanki-app sets none: app/src/lower.rs lowers every user with speed_limit: None. katana sets each user’s limit from the node’s rate and the user’s rate, the smaller of the two that are not zero (src/lower/mod.rs → determine_rate); the katana lowering page and the speed limits guide describe where those rates come from.
The bucket
Section titled “The bucket”PaceState::refill(now) credits the time since last at the current rate, up to one second of rate, and moves last to now. An unlimited state only moves last. The elapsed time is now.saturating_duration_since(last), so a last later than now credits nothing.
Pacer::set_rate(rate) does nothing when the rate is unchanged, so publishing the same limits on every apply leaves every bucket as it is. Otherwise it refills at the old rate, then sets the balance:
| Change | New balance | Why |
|---|---|---|
| Limited to unlimited | 0 | Any debt is forgiven |
| Unlimited to limited | new (a full bucket) |
A newly limited user starts with a full burst |
| Limited to another limit | min(balance, new) |
Credit earned at the old rate is kept up to the new burst; debt is kept |
It then stores the new rate in limited with Release ordering.
Pacer::charge(n) takes the n bytes a read or write has moved. When limited is 0 it returns after one Acquire load, without the lock. Otherwise it locks and checks rate again under the lock, returning if it is 0 (a set_rate to unlimited may have run between the load and the lock); then it refills and subtracts n, even past zero.
Pacer::poll_ready(cx, pause) decides whether the next read or write may start:
- When
limitedis 0, clearpauseand returnReady. - Otherwise lock, refill, and return
Ready(clearingpause) when the rate is 0 or the balance is not negative. - Otherwise compute the wait as
-tokens / rateseconds, clamped to betweenMIN_WAIT= 1 ms andMAX_WAIT= 1 s, and release the lock. Setpause’sSleepto the deadlinenow + wait: the first time it allocates one withsleep_until, and on every pass it callsreset(deadline)on it. Then poll it. If it is pending, returnPending; if it already elapsed, loop back to step 2.
MAX_WAIT bounds each sleep so a raised or removed limit is felt within one second (a_raised_or_removed_limit_is_felt_within_a_second); a flow in debt wakes at least once a second and looks again. MIN_WAIT keeps a sliver of debt from spinning the timer.
A paced poll
Section titled “A paced poll”flowchart TB
poll["poll_read, poll_write, poll_send_to or poll_recv_from"]
k1{"killed?"}
p{"flow has a pacer?"}
ready{"Pacer::poll_ready"}
park["register kill waker, check kill again, Pending"]
op["poll Guarded, then the link"]
moved{"Ready Ok?"}
add["add_up or add_down: counter, wire, Pacer::charge"]
err["Err: the flow was killed"]
out["return the result"]
poll --> k1
k1 -->|yes| err
k1 -->|no| p
p -->|no| op
p -->|yes| ready
ready -->|"Pending, in debt"| park
ready -->|Ready| op
op --> moved
moved -->|yes| add --> out
moved -->|no| out
The wait comes before the operation and the charge after it. So a read that is held back leaves the bytes in the outbound (a TCP destination’s sends back up in its window, and a UDP link’s packets wait in its socket), and a write that is held back returns Pending, which keeps the runtime’s forward queued and so stops it reading the client (The server runtime). The pacing timer and the kill waker are both registered with the same task waker, so either wakes the parked poll.
Why a transfer runs into debt
Section titled “Why a transfer runs into debt”Neither the runtime nor the fan-out moves bytes in fixed sizes: a read is as large as the staging room allows, and a write is whatever range the protocol core forwarded, both bounded by per-protocol buffer sizes. A limiter that waited until the bucket held n tokens before moving n could never let a chunk larger than the one-second burst through without exceeding the rate, and would give different speeds to different protocols.
The pacer instead:
- lets a transfer start whenever the balance is not negative;
- charges its full size after it moved, even past zero;
- lets none of the user’s flows start another transfer until the refill has paid the debt back.
Every byte is paid for at the rate, whether it moved in one piece or in many. The only slack is the burst (at most one second of rate, banked while idle) and the transfers already under way when the balance crossed zero; the transfers after them pay for them. Worked example, from a_raised_or_removed_limit_is_felt_within_a_second, at 1,000 bytes per second:
| Step | Balance before | Moves | Balance after | Next transfer |
|---|---|---|---|---|
| The limit is set (unlimited to 1,000 B/s) | 0 | nothing | 1,000 | at once |
| One write of 11,000 bytes | 1,000 | all 11,000 bytes | −10,000 | in 10 s, looked at in sleeps of at most 1 s |
| The limit is removed 0.1 s later | −9,900 | nothing | 0 | at the end of the current sleep, within 1 s |
With the limit kept, the same arithmetic gives a_limited_users_flow_moves_at_its_limit: 300,000 bytes plus one at 100,000 B/s take between 2 and 3 seconds, one second’s burst and then 200,000 bytes of debt.
One bucket per user
Section titled “One bucket per user”add_up and add_down both charge the same pacer, and every flow of a user holds the same Arc<Pacer>, whatever its inbound, session or direction. Two flows of one user, one uploading 150,000 bytes and one downloading 150,000 bytes at 100,000 B/s, share one burst and one debt, so one more byte moves only after 2 seconds (a_users_flows_share_one_limit). A flow of an unlimited user takes the fast path, one atomic load of limited in poll_ready and one in charge; an anonymous flow (an open inbound, a shared credential, a TUN device) has no pacer at all (an_unlimited_user_and_an_anonymous_flow_are_not_paced). A limited user’s flows take the pacer’s mutex twice per read or write, once in poll_ready and once in charge.
The registry and why scc
Section titled “The registry and why scc”The registry is written on every flow open and close, from every worker thread, and scanned in full by every sampler tick, flows() and kill_where. It has to keep opens and closes fast while a scan is running. In supervisor/Cargo.toml one comment covers both scc and futures: “The live-flow registry, and the waker a flow is killed through.” Cargo.lock resolves scc to 3.8.8.
The commit that introduced the tracker records the comparison that chose it: 8 threads churning opens and closes over 500,000 live entries, on 32 cores.
| Map | Open and close throughput, idle | Beside full scans |
|---|---|---|
scc::HashMap |
105 M ops/s | 62 M ops/s (a scan took 5.4 ms) |
dashmap |
86 M ops/s | 10 M ops/s |
papaya |
11 M ops/s | not recorded |
Mutex<HashMap> |
5 M ops/s | 15 k ops/s |
That comparison is not part of the repository’s benchmarks.
The benchmark
Section titled “The benchmark”supervisor/benches/tracking.rs is a criterion benchmark. supervisor/Cargo.toml declares it as [[bench]] name = "tracking" with harness = false, and criterion is a dev-dependency. Run it with:
cargo bench -p etemenanki-supervisor --bench trackingIt has three groups. The two usage groups (take_usage, 50k active users and reconcile (sampler tick), 50k active users) are described on Per-user usage accounting. The third measures what a metered flow adds per byte:
- Group
relay 256 MiB in 16 KiB reads, withThroughput::Bytes(256 << 20)andsample_size(20), on a current-thread Tokio runtime. untrackeddrainstokio::io::repeat(0).take(256 MiB)in 16 KiB reads, as a runtime reads an outbound.metereddrains the same reader wrapped bytrack::bench::metered, a#[doc(hidden)]public helper (“What the benches need of the internals; not part of the API”) that registers it as one flow of a tracker of its own. The helper builds a freshSessions(a newCancellationToken, a newTaskTrackerand a default ledger), callsTracker::newand drops the sampler, and meters the stream throughmeter_streamwith a fixedFlowMeta: no session, inboundbench,Principal::anonymous(), no source, TCP destinationbench:443, not sniffed, outboundbench@v1, no rule, plane epoch 0.
The metered flow has no pacer and no subscriber, so the difference is the kill check and the counter add per read. When the benchmark was added it measured 153.3 GiB/s untracked against 153.1 GiB/s metered (release build, 32 cores): the per-poll work is lost in the cost of a 16 KiB copy.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| A flow is in the registry exactly while its wrapper lives | open inserts before the wrapper exists; FlowHandle::drop removes |
a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped, killing_a_flow_fails_its_pending_read |
| Only an opened outbound is a flow | meter_stream runs after connect_stream resolved Ok; meter_datagram after the sub-link opened |
By construction |
| A flow counts exactly the payload that moved | add_up/add_down only after Ready(Ok), with the inner result’s size |
a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped (5 up, 11 down), mux_sub_flows_and_their_session_count_exactly_what_moved, a_udp_association_counts_each_sub_link_and_charges_its_session |
| Each mux sub-flow and each UDP sub-link is a flow of its own, under the connection’s session | The runtime calls connect per key; the fan-out meters per sub-link with its FlowScope |
mux_sub_flows_and_their_session_count_exactly_what_moved, a_udp_association_counts_each_sub_link_and_charges_its_session (two outbounds, two flows, one session) |
Opened precedes Closed, and Closed carries the final counts |
open sends Opened after insertion; drop sends Closed after removal |
a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped |
| Opens and closes never wait for a subscriber | Bounded broadcast; send errors ignored |
None |
| A kill fails the flow’s next poll in either direction, and wakes a parked one | Flag checked before and after registering the waker; kill wakes both wakers |
killing_a_flow_fails_its_pending_read, killing_a_paced_flow_fails_its_pending_write |
| A kill ends one flow only | Other flows have their own flags; cores close only the failed key | a_selective_kill_leaves_the_other_flows, killing_one_mux_sub_flow_leaves_its_siblings_and_the_carrier |
kill_where counts each flow once |
Skips killed() entries |
a_selective_kill_leaves_the_other_flows |
| Every byte is counted in exactly one tick, live or closed | Sampler-only watermarks; scan before drain | the_sampler_counts_every_byte_once_by_inbound_outbound_and_user |
Rates are the tick’s bytes over the interval passed to tick |
per_second divides by interval |
the_sampler_counts_every_byte_once_by_inbound_outbound_and_user (50 bytes over 0.5 s is 100 B/s) |
| That interval is the time measured since the previous tick | run passes now - last on Tokio’s clock |
None; the tests call tick with explicit intervals |
A user with nothing live and no traffic in the tick is left out of users |
Group::active filter |
the_sampler_counts_every_byte_once_by_inbound_outbound_and_user |
| One publish per interval, the first after one interval, and one more at stop | interval_at(now + interval, …); one more tick when stop wins the select! |
the_sampler_publishes_once_per_interval_and_once_when_stopped |
| The last snapshot holds every flow’s last bytes | The sampler stops only after every connection ended | mux_sub_flows_and_their_session_count_exactly_what_moved checks stats while live; shutdown order has no dedicated test |
| A limit is a property of bytes, not of chunk sizes | Charge after the move, into debt; wait before the next | a_limited_users_flow_moves_at_its_limit |
| One bucket per user, both directions, all flows | One Arc<Pacer> per user key; both add_* charge it |
a_users_flows_share_one_limit |
| Unlimited and anonymous flows are never delayed | limited == 0 fast path; no pacer without a user key |
an_unlimited_user_and_an_anonymous_flow_are_not_paced |
| A raised or removed limit is felt within 1 s | MAX_WAIT caps every pacing sleep |
a_raised_or_removed_limit_is_felt_within_a_second |
| Changing a limit keeps the user’s sessions | The limit is not part of admission | a_speed_limit_change_keeps_the_users_sessions (supervisor/tests/hot_swap.rs) |
| A refused apply or user edit leaves the live limits alone | publish_speed_limits runs only in the commit phase |
None |
| Re-publishing an unchanged limit leaves the bucket as it is | set_rate returns early on an equal rate |
None |
| A user in several sets gets the smallest limit | and_modify(min) in publish_speed_limits |
None |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Situation | What happens |
|---|---|
| The dial of a TCP flow fails | No flow is registered; the runtime hands the core ConnectFailed |
| A UDP sub-link fails to open | No flow is registered and nothing is counted; poll_opening logs udp fan-out: opening <tag>@v<version> failed: <error> at debug (Outbounds, UDP fan-out and balancers) |
The session refuses the flow (Session::admit) |
connect fails before routing; nothing is registered (Users, principals and sessions) |
| A flow is killed | Every later poll fails with ConnectionAborted, the flow was killed; see What the connection does with the error |
| A drain policy closes the flow’s outbound version | Guarded, inside the wrapper, fails the flow with the outbound this flow was opened on was drained (The plane); the runtime or the fan-out handles the error as it handles a kill, and the flow deregisters when its wrapper drops |
| Nobody subscribes | Event sends fail and are ignored |
| A subscriber falls more than 1024 events behind | Its next recv returns RecvError::Lagged(n); it resumes from the oldest event kept |
| The sampler has stopped | Closed-queue sends fail and are ignored; stats() receivers keep the last snapshot, and changed() reports the sender gone |
SupervisorBuilder::sample_interval(Duration::ZERO) |
start panics with the sample interval must not be zero before it applies anything |
| A panic inside a tick | Re-raised in the sampler task with resume_unwind; the task ends and no further snapshots are published |
Cancellation is by drop, like everything on the data plane. The wrappers spawn nothing: dropping one drops its link, its pacing timers and its FlowHandle, which deregisters the flow. The sampler stops only through background_stop, which the actor cancels at the end of its shutdown, or when a tick panics (see the table).
Limits
Section titled “Limits”| Constant or bound | Value | Defined in | Meaning |
|---|---|---|---|
EVENT_CAPACITY |
1024 | supervisor/src/track/mod.rs |
Events a subscriber may fall behind before it loses them |
DEFAULT_SAMPLE_INTERVAL |
1 s | supervisor/src/track/sampler.rs |
The tick when sample_interval is not set; every front end keeps it |
MAX_WAIT |
1 s | supervisor/src/track/pace.rs |
Longest pacing sleep before a flow looks at its bucket again |
MIN_WAIT |
1 ms | supervisor/src/track/pace.rs |
Shortest pacing sleep |
| Burst | One second of the user’s rate | supervisor/src/track/pace.rs → PaceState::refill |
The most a user banks while idle |
MAX_SUBS |
64 | supervisor/src/topology/outbound/udp_fanout.rs |
Sub-links, and so flows, one UDP association keeps; the least recently sent-to goes first |
| Closed queue | Unbounded mpsc, emptied by every tick |
supervisor/src/track/mod.rs |
One Arc<FlowEntry> per flow closed since the last tick |
Each live flow costs one registry entry and one Arc<FlowEntry> (its FlowMeta, four AtomicU64, one AtomicBool and two AtomicWaker), plus the wrapper’s FlowHandle (up to three more Arcs) and, only while its user is in debt, one boxed Sleep per direction. Nothing is done per byte; per read or write, an unlimited flow adds two Acquire loads of the kill flag (three when the link returns Pending: paced checks, guarded checks, and guarded checks again after registering the waker), two loads of the pacer’s rate when it has a pacer, and one relaxed atomic add, plus one more for a charged Wire.
Unit tests in supervisor/tests/unit/track.rs, compiled into the crate as track::tests. They build a tracker over a fresh Sessions and meter tokio::io::duplex streams with hand-made FlowMeta. The sampler tests drive Sampler::tick directly with explicit intervals, or spawn Sampler::run on a paused clock. The sampler-run test and the pacing tests use #[tokio::test(start_paused = true)], which needs Tokio’s test-util feature; supervisor/Cargo.toml enables it on its tokio dev-dependency.
| Test | Behaviour it pins |
|---|---|
a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped |
Opened on registration; 5 bytes up and 11 down counted; the flow leaves the registry on drop; Closed carries the final counts |
killing_a_flow_fails_its_pending_read |
A kill wakes a read parked on a link that never answers, with ConnectionAborted; kill returns false once the wrapper is gone |
a_selective_kill_leaves_the_other_flows |
kill_where by outbound tag kills two of three flows, counts 0 the second time, and the third flow still writes |
the_sampler_counts_every_byte_once_by_inbound_outbound_and_user |
Totals, rates and live counts per total, inbound, outbound and user over three ticks; a closed flow’s last bytes counted once; an idle user left out |
the_sampler_publishes_once_per_interval_and_once_when_stopped |
With run spawned on the paused clock: ticks at 1 s and 2 s, and tick 3 at the stop |
a_limited_users_flow_moves_at_its_limit |
300,000 bytes plus one at 100,000 B/s take at least 2 s and less than 3 s |
an_unlimited_user_and_an_anonymous_flow_are_not_paced |
Another user and an anonymous flow move 300,000 bytes with no delay while user 1 is limited |
a_users_flows_share_one_limit |
An upload and a download of one user share one burst and one debt |
a_raised_or_removed_limit_is_felt_within_a_second |
A flow ten seconds in debt resumes within 1 s of the limit’s removal |
killing_a_paced_flow_fails_its_pending_write |
A kill ends a write parked on its user’s debt in under a second |
Integration tests in supervisor/tests/tracking.rs run a real supervisor on loopback sockets. They wait on the supervisor with two helpers: eventually polls a probe every 10 ms until it returns a value and fails with the condition never held after SOON = 5 s, and sessions_gone uses it to wait until tracker().sessions() is empty.
| Test | Behaviour it pins |
|---|---|
mux_sub_flows_and_their_session_count_exactly_what_moved |
Three VLESS mux sub-flows are three flows of one session, each with its inbound, label, outbound direct, no rule, and exactly its payload; the snapshot’s user row reaches the payload sum (sampled every 50 ms); the flows are gone after the connection |
killing_one_mux_sub_flow_leaves_its_siblings_and_the_carrier |
The killed sub-flow ends; its sibling still echoes; a new sub-flow opens; the session count stays 1 |
a_udp_association_counts_each_sub_link_and_charges_its_session |
A SOCKS association routed to two outbounds is two UDP flows, one per outbound, with 10 and 5 bytes each way, rule None and rule 0, and one session |
a_hysteria2_session_bills_its_quic_connections_bytes |
The Hysteria 2 stream’s flow counts exactly its 20,000 bytes each way while the session counts more |
The same file’s usage-sink tests and a_hysteria2_user_set_admits_by_password are described on Per-user usage accounting and Users, principals and sessions. supervisor/tests/hot_swap.rs → a_speed_limit_change_keeps_the_users_sessions sets a limit through set_users on a live mux session and checks that the session id is unchanged.
webclient/tests/tracking.rs exercises the tracker end to end through the REST router, over a supervisor with one SOCKS inbound and a direct outbound, sending requests with oneshot:
| Test | Behaviour it pins |
|---|---|
connections_lists_live_sessions_and_flows |
One session and one flow; the session’s wire bytes include the SOCKS greeting and request while the flow counts the 5-byte payload each way; the flow names its session, tcp, the inbound, direct and no rule |
deleting_a_flow_closes_it_alone |
DELETE /v1/connections/{flow} answers 204 and closes that connection; the other still echoes; a second delete of the same id answers 404 |
deleting_by_selector_closes_the_matching_flows |
A selector matching nothing closes 0; an unknown query key answers 400; outbound=direct&inbound=socks-in closes 2; no selector closes every flow |
traffic_streams_one_event_per_tick |
At a 100 ms interval, /v1/traffic sends one traffic event per tick with consecutive tick numbers, and the total reaches the 1,000 bytes sent |
connection_events_stream_opens_and_closes |
/v1/connections/events sends opened, then closed for the same id with 7 bytes each way |
Run them with cargo test -p etemenanki-supervisor --lib track, cargo test -p etemenanki-supervisor --test tracking and cargo test -p etemenanki-webclient --test tracking.
Nothing tests a subscriber’s lag, a killed UDP sub-link, publishing limits only at commit, the smallest limit across sets, an unchanged limit leaving the bucket alone, MIN_WAIT, or the sampler’s behaviour when a tick is late. Add one when you touch those paths. The harnesses are described on Testing.