Per-user usage accounting
Source files: 29 · checked against Etemenanki 555b7df · katana v4.1.1
Etemenanki/supervisor/src/entity/usage.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/entity/user.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/serve.rsEtemenanki/supervisor/src/connector.rsEtemenanki/supervisor/src/build/users.rsEtemenanki/supervisor/src/build/inbound.rsEtemenanki/supervisor/src/topology/flow.rsEtemenanki/supervisor/src/topology/spec_plan/inbound.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/metered.rsEtemenanki/supervisor/src/track/sampler.rsEtemenanki/supervisor/benches/tracking.rsEtemenanki/supervisor/Cargo.tomlEtemenanki/concepts/src/runtime.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/hysteria/obfs.rsEtemenanki/app/src/instance.rsEtemenanki/ffi/src/proxy.rsEtemenanki/webclient/src/lib.rsEtemenanki/webclient/src/tracking.rsEtemenanki/supervisor/tests/unit/usage.rsEtemenanki/supervisor/tests/unit/session.rsEtemenanki/supervisor/tests/tracking.rskatana/src/manager/node.rs
The supervisor counts, for every user, the bytes their client’s connections put on the server’s wire, and hands those counts to a front end as deltas: by polling (Supervisor::take_usage) or by pushing to a UsageSink on an interval. Each byte a bound session counts, up to its close, lands in exactly one delta, whichever of the two took it. The counting lives in each session’s Wire; the accounting lives in one Ledger per supervisor, which keeps one account per user with a watermark per live session.
This page is for contributors who change supervisor/src/entity/usage.rs, and for front-end authors who bill from the supervisor. It covers what is counted and what is not, the Wire, the Ledger and its accounts, the exactly-once mechanism, the pull and push paths including the push task at shutdown, the cost, and how usage differs from the payload figures in StatsSnapshot. Sessions themselves (how they open, bind a principal and end) are on Users, principals and sessions. Flows, the sampler and speed limits are on Tracking.
Responsibilities
Section titled “Responsibilities”The usage code:
- gives every session a
Wirethat counts the bytes its client’s side carried, from the session’s first byte; - opens the user’s account in the
Ledgerwhen the session’s first flow binds it to a user, and folds the session’s remainder into that account when the session ends; - moves each live session’s new bytes into its user’s pending total once per sampler tick, off the async workers;
- hands out pending totals as
UsageDeltavalues, at most one per user per take, and never a byte twice; - with a
UsageSinkinstalled, runs a push task that reconciles, takes and reports at each tick of its interval, and once more after shutdown has waited for every connection to end; - reports each user’s running total since the supervisor started (
usage_snapshot) without taking anything.
It does not:
- count payload. The payload each flow moves is the sampler’s business, and reaches
StatsSnapshot(see Usage andStatsSnapshot). - deliver anything to a panel. A delta handed out is the front end’s to deliver; the supervisor never offers it again.
- persist. The ledger lives in memory, and
Actor::newbuilds an empty one for every supervisor. - pace anyone. Speed limits charge payload per flow, through the pacer on Tracking.
Who reads usage
Section titled “Who reads usage”| Front end | Path | Where |
|---|---|---|
| katana | Pull: calls Supervisor::take_usage from its report cycle. It starts each node’s supervisor with Supervisor::start(spec), so the sampler ticks at the default 1 s and no sink is installed |
src/manager/node.rs → report_traffic; see katana traffic accounting |
| etemenanki-app | Neither: Instance::start builds its supervisor with Supervisor::builder() and installs no sink |
app/src/instance.rs |
etemenanki-ffi |
Neither: it builds its supervisor with Supervisor::builder().socket_options(...) and installs no sink |
ffi/src/proxy.rs |
REST API (etemenanki-webclient) |
Neither: its traffic view leaves per-user figures out, and per-user usage “is the supervisor’s Rust API alone” | webclient/src/lib.rs, webclient/src/tracking.rs |
The REST API itself is described on its own page, The REST API (etemenanki-webclient).
What is counted
Section titled “What is counted”Usage is wire bytes per session: what the server’s bandwidth carries for the user’s client, protocol overhead included. The module docs of supervisor/src/entity/usage.rs define it, and each kind of session feeds its Wire from a different place.
| Session | Wire |
Fed by | What the count covers |
|---|---|---|---|
| A stream inbound’s connection (HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022) over plain TCP, TLS or WebSocket | Wire::counted() of the socket’s session, which run_stream_inbound opens at accept, before the transport’s own handshake |
serve.rs → drive: every Traffic step the runtime yields adds transport_rx as up and transport_tx as down |
The protocol stream after its transport and before the protocol’s framing and encryption are undone: the protocol’s handshake, headers, padding and ciphertext are counted. The TLS or WebSocket handshake, TLS records and WebSocket frames are not: nothing adds to the wire until drive runs on the stream the transport yields |
| A stream on a gRPC carrier | Each HTTP/2 stream is a session of its own, with its own Wire::counted() |
drive, as above |
The stream’s bytes after the gRPC framing |
| A SOCKS connection | Wire::counted() of the socket’s session |
serve.rs → WireCounted, a wrapper around the client stream: each successful poll_read adds up, each successful poll_write adds down |
Every byte of the control stream: greeting, authentication, request and replies, and a CONNECT relay’s bytes |
| A SOCKS UDP association | The same session’s Wire |
AppConnector::charging_datagrams → FlowScope::charge → each sub-link’s MeteredDatagram; FlowHandle::add_up and add_down also add to the charged wire |
The payload of every datagram relayed to and from a destination. The datagrams never cross the control stream, so this is the only way they reach the session |
| A Hysteria 2 QUIC connection | Wire::quic(ConnectionBytes), built in run_hysteria_inbound |
quinn::Connection::stats(): udp_rx.bytes as up and udp_tx.bytes as down |
The UDP payload of every datagram quinn receives or sends on the connection: the QUIC handshake, QUIC framing, TLS and the HTTP/3 credential exchange included. IP and UDP headers are not counted. With obfuscation, neither is the salt of SALT_LEN (8) bytes that the Salamander layer (protocols/src/hysteria/obfs.rs) adds to every datagram below quinn, where quinn never sees it |
| A Unix-socket listener’s connection | Wire::counted() of the socket’s session, which serve_socket serves with no source address (connector_for(tag, None, carrier)) |
drive or WireCounted, as for TCP |
The socket’s bytes; a Unix listener carries no transport |
Directions are the user’s: up is what their client sent toward destinations, down is what came back to it.
Hysteria 2 is therefore the one session kind whose count includes its transport’s own overhead: the QUIC connection is both the transport and the session. The test a_hysteria2_session_bills_its_quic_connections_bytes pins it: a 20,000-byte echo shows exactly 20,000 bytes each way on the flow and more than 20,000 each way on the session and in the user’s usage.
When a session starts counting toward a user
Section titled “When a session starts counting toward a user”A session counts toward its user once its first flow binds it. Session::admit binds the session to the Principal of the first flow it admits; if that principal carries a UserKey, the admit calls Ledger::bind with the session’s Wire. The binding starts the session’s watermark at (0, 0), so everything the wire counted before the bind is reported with the rest: a stream session’s protocol handshake, or a Hysteria 2 connection’s QUIC handshake and credential exchange.
AppConnector::connect runs that admit for every flow before it routes or dials anything (see Routing and the plane). The bind therefore does not depend on the dial: a session whose first flow fails to open is bound all the same, and what its wire counted, the protocol handshake included, is reported under its user.
A later flow presenting another principal does not rebind the session (the_first_flow_binds_the_session_to_its_principal). A mux connection is one session, so every sub-flow’s bytes count toward the user who opened the first.
What counts toward nobody
Section titled “What counts toward nobody”| Case | Why nothing reaches an account |
|---|---|
| A session that never opens a flow: a transport or protocol handshake that fails or times out, a Hysteria 2 client that never authenticates, or one that authenticates and never opens a stream or a datagram flow | Ledger::bind runs only from the first admitted flow. Until a core is established (its request parsed and its flow open or opening), drive waits at most HANDSHAKE_TIMEOUT (10 s, protocols/src/core/mod.rs) for each runtime step and otherwise ends the connection with inbound handshake timed out after 10s; such a connection has opened no flow, so what its wire counted goes to nobody |
| A session closed before its first flow | Session::admit refuses every flow with the session was closed (PermissionDenied) once the session’s token is cancelled or its registry entry is gone, before it looks at the principal |
| A session whose first flow presents a principal revoked before the bind | Session::admit refuses it with the user was removed before the session opened a flow and returns before it reaches Ledger::bind |
| A flow admitted as nobody: an inbound in its open or shared mode | Its principal is Principal::anonymous(), whose user_key() is None. build/inbound.rs → user_table gives it to an inbound without a user set (InboundSpec::users, defined in topology/spec_plan/inbound.rs, is None): SOCKS without auth, HTTP without accounts, Shadowsocks with only the server password, Shadowsocks 2022 with only the server PSK, and Hysteria 2 with a shared_password |
| A gRPC carrier’s own framing | Each HTTP/2 stream on a gRPC carrier is a session of its own and counts its bytes after the gRPC framing; the carrier’s framing counts toward nobody |
| A TUN device | A device admits no users, so its connector has no session (run_tun_inbound) |
| Flows the supervisor opens itself | They carry anonymous_user() (topology/flow.rs) |
Payload per user, including anonymous traffic by inbound and outbound, is in the sampled StatsSnapshot instead.
Key types
Section titled “Key types”UsageDelta and UsageTotal
Section titled “UsageDelta and UsageTotal”#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]pub struct UsageDelta<U: UserId> { pub user: U, /// Bytes from the user's client, toward destinations. pub up: u64, /// Bytes back to the user's client. pub down: u64,}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]pub struct UsageTotal<U: UserId> { pub user: U, pub up: u64, pub down: u64,}UsageDeltais traffic one user moved since the last take.take_usageand a sink batch carry at most one per user, and none whoseupanddownare both zero.UsageTotalis everything one user moved since the supervisor started, reported or not.- Both are keyed by the front end’s own
U: UserId, not by the internalUserKey. Both derive serde’sSerializeandDeserialize, which apply wheneverUimplements them, so a front end can send a report as it is; that is whysupervisor/Cargo.tomldepends onserde.
pub trait UserId: Clone + Eq + Ord + Hash + Debug + Send + Sync + 'static { /// The name nodes surface for this user. fn name(&self) -> UserName;}
impl UserId for UserName { /* returns a clone of itself */ }U is whatever the front end already calls its users: a panel’s numeric id, or UserName (an Email or a Username) itself, which implements the trait. The model behind it is on Users, principals and sessions.
UsageSink
Section titled “UsageSink”pub trait UsageSink<U: UserId>: Send + Sync + 'static { /// Record one batch: at most one delta per user, none for a user who /// moved nothing. Never called with an empty batch. fn report(&self, batch: &[UsageDelta<U>]) -> impl Future<Output = ()> + Send;}reportruns on the push task, never on the data plane. A slow sink delays only its own next batch, which then covers the longer interval.- A delta handed to a sink is not also returned by
take_usage. reportreturns(). The supervisor cannot learn that a batch was not delivered, and it never hands a batch out again: a sink that must not lose one keeps it until it is delivered.
pub struct Wire { up: AtomicU64, down: AtomicU64, /// Read from the connection instead of counted, for a QUIC session. quic: Option<ConnectionBytes>,}
impl Wire { pub fn counted() -> Self; pub fn quic(bytes: ConnectionBytes) -> Self; pub fn add(&self, up: u64, down: u64); pub fn read(&self) -> (u64, u64);}| Method | Behaviour |
|---|---|
counted() |
Both counters at zero, no QUIC reader. |
quic(bytes) |
The same counters plus bytes, the reader of one QUIC connection’s byte counts. |
add(up, down) |
A fetch_add with Ordering::Relaxed per direction, skipped for a zero. This is the whole per-byte cost on the data plane. |
read() |
(up, down) so far, and never decreasing. For a QUIC wire it returns the connection’s (received, sent) plus the counters. Nothing adds to a Hysteria 2 session’s counters, so its read is the connection’s counts. |
Relaxed ordering is enough because nothing is published through the counters: each reader only needs a value that never goes backwards, and the ledger’s watermarks are protected by the account’s lock, not by the atomics.
ConnectionBytes
Section titled “ConnectionBytes”#[derive(Clone)]pub struct ConnectionBytes(quinn::Connection);
impl ConnectionBytes { /// `(received, sent)` so far, from the server's side. Never decreases. pub fn bytes(&self) -> (u64, u64);}Hy2Inbound::run builds one from each connection right after its QUIC handshake completes and passes it with the client’s address to the make_connector callback, which is where run_hysteria_inbound opens the session with Wire::quic(bytes). bytes() reads stats().udp_rx.bytes and stats().udp_tx.bytes, which takes the connection’s internal lock. ConnectionBytes holds a clone of the connection, so it keeps counting for as long as it is held; the Session, its registry entry and its ledger record hold it through the Wire (see Who holds a wire), and let go of it when the session ends.
Ledger and Account
Section titled “Ledger and Account”#[derive(Default)]pub struct Ledger { inner: Mutex<LedgerInner>,}
#[derive(Default)]struct LedgerInner { /// Every user that has bound a session. accounts: HashMap<UserKey, Arc<Account>>, /// Those with a live session or bytes not yet reported: what a take visits. active: HashMap<UserKey, Arc<Account>>,}
pub struct Account { state: Mutex<AccountState>,}
#[derive(Default)]struct AccountState { live: HashMap<SessionId, LiveSession>, /// Bytes reconciled or folded in, not yet taken. pending: (u64, u64), /// Everything taken so far. reported: (u64, u64),}
struct LiveSession { wire: Arc<Wire>, /// What of this session was moved into `pending` so far. settled: (u64, u64),}Both mutexes are parking_lot::Mutex, never held across an .await. Ledger, Account and Wire and their methods are pub in etemenanki_supervisor::entity::usage, not re-exported at the crate root. No front end uses them; the benchmark drives Ledger::bind, take and reconcile directly. entity/user.rs → Account (a user name and password pair) is a different type with the same name.
The SessionId a live session is filed under comes from Sessions::next, an AtomicU64 counting from 1, and prints as session#<n>; it is unique for the supervisor’s lifetime.
| Method | Locks | What it does |
|---|---|---|
Ledger::bind(user, session, wire) -> Arc<Account> |
Ledger, then the account | Finds or creates the user’s account in accounts, inserts session into its live map with settled = (0, 0), and inserts the account into active. Returns the account, which the Session keeps in a OnceLock to close later. Adding the session under the ledger’s lock means a concurrent take either finds the account active or finds it with this session. |
Ledger::reconcile() |
Ledger to list active, then each account in turn |
Clones the active accounts into a Vec under the ledger’s lock and releases it. For each account, under its lock, settles every live session into pending. |
Ledger::take() -> Vec<(UserKey, u64, u64)> |
Ledger, and each active account in turn | For each active account: moves pending into reported and resets it, pushes (user, up, down) when that was not (0, 0), and keeps the account in active only if it still has a live session. |
Ledger::totals() -> Vec<(UserKey, u64, u64)> |
Ledger, and each account in turn | For every account ever bound: reported + pending + (wire.read() - settled) over its live sessions. Marks nothing. |
Account::close(session) |
The account | Removes the session from live and adds what it moved past its watermark to pending. |
LiveSession::settle is the one place bytes leave a wire for the ledger:
fn settle(&mut self) -> (u64, u64) { let now = self.wire.read(); let due = (now.0 - self.settled.0, now.1 - self.settled.1); self.settled = now; due}Neither take nor totals promises an order: both walk a HashMap.
The supervisor’s side: UsageBook
Section titled “The supervisor’s side: UsageBook”struct UsageBook<U> { ledger: Arc<Ledger>, /// Every id a key was handed out for, in key order; appended to as the /// actor commits new keys, never shrunk. names: RwLock<Vec<U>>,}
impl<U: UserId> Supervisor<U> { pub fn take_usage(&self) -> Vec<UsageDelta<U>>; pub fn usage_snapshot(&self) -> Vec<UsageTotal<U>>;}
impl<U: UserId> SupervisorBuilder<U> { pub fn usage_sink(self, sink: impl UsageSink<U>, interval: Duration) -> Self;}- The ledger speaks
UserKey; a front end speaksU.UsageBooktranslates withnames, aparking_lot::RwLock<Vec<U>>in which keynnamesnames[n - 1](build/users.rs→id_of).UserKeyshands keys out from 1 on first sight of an id and keeps each id’s key, so a removed user’s last bytes are still reported under their id, and a user added back gets the same key and the same account. - A key is handed out only when an inbound admits the user:
build/users.rs→admitcallsUserKeys::keyfor a user of the inbound’s set who holds the credential kind the protocol reads. Being in a set is not enough, so a user no inbound admits never gets a key, never binds a session and never gets an account. Actor::commit_keysappends the ids of newly handed-out keys tonamesat the start of every commit, before any user table is stored or listener started, in an apply and in a user edit alike. Keys are handed out on a staged copy during prepare and kept only if the change commits. So every key is published before a session can bind it.UsageBook::namelooks a key up and panics with<key> was bound before it was publishedwhen it finds none.UserKeyhas noDisplay, so the message formats the key with{key:?}, for exampleUserKey(7) was bound before it was published. Taken bytes cannot be put back, so a key without a name is treated as a bug, not as a delta to drop.UsageBook::takecallsLedger::takeand maps each tuple to aUsageDelta<U>under the read lock;UsageBook::totalsdoes the same withLedger::totalsandUsageTotal<U>.take_usageandusage_snapshotare plain synchronous methods on the handle. They read the book directly and never go through the actor’s command channel, so they answer while an apply is in progress and aftershutdownhas returned. They run on the caller’s thread and take the ledger’s and each account’sparking_lotlocks, so an async caller holds its worker for the whole walk.- The
Actorowns anArc<UsageBook<U>>;startclones it into theSupervisorhandle and into the push task.
Where the data plane feeds a wire
Section titled “Where the data plane feeds a wire”// The session's wire is fed from the runtime's own counts: what the// transport side moved is what the client's side carried.let wire = connector.session().map(|s| s.wire().clone());let count = |traffic: &Traffic| { if let Some(wire) = &wire { wire.add(traffic.transport_rx, traffic.transport_tx); }};drive runs the runtime in showing_progress mode, so each Traffic it yields is a delta, and the deltas sum to the runtime’s total (The server runtime). count runs for handshake steps and relay steps alike. A step that ends in an error yields no Traffic: whatever it had read or written before failing stays in the runtime’s unreported delta (poll_once hands the delta out only on success), so the wire never sees it. drive converts that error with io::Error::other(e.to_string()) and returns it, and serve_connection logs it at debug as <protocol> connection from <source> ended: <error>, the source in Debug form (Some(<ip>) or None).
struct WireCounted<S> { inner: S, wire: Option<Arc<Wire>>,}The SOCKS driver (SocksInbound::serve, see SOCKS) runs no ProxyServerRuntime, so nothing reports what its stream moved. serve_connection wraps the stream in WireCounted instead: poll_read adds the bytes the read filled as up, and poll_write adds the count it returned as down. Flush and shutdown pass through. The same call hands the driver connector.charging_datagrams(), which sets AppConnector::charge_datagrams. When the association’s FanOutLink opens a sub-link, poll_opening passes FlowScope::charge (the session’s Arc<Wire>) to Tracker::meter_datagram, and the sub-link’s FlowHandle adds every payload byte it sends or receives to that wire as well as to the flow’s own counters. See Outbounds for the fan-out.
Who holds a wire
Section titled “Who holds a wire”A session’s Wire is created once, by whoever opens the session, and moved into Sessions::open, which wraps it in an Arc. Three long-lived holders then keep it:
| Holder | Field | Released |
|---|---|---|
The Session handle |
Session::wire, returned by Session::wire() to drive, to serve_connection for SOCKS and to AppConnector::connect for a charged association |
With the last Arc<Session> |
| The registry entry | Entry::wire in Sessions, read by Sessions::stats |
Sessions::remove, from Session’s Drop |
| The user’s account | LiveSession::wire, from Ledger::bind |
Account::close, from Session’s Drop |
The per-connection feeders hold clones for as long as they feed it: drive’s count closure, WireCounted::wire for SOCKS, and for a SOCKS UDP association FlowScope::charge in the FanOutLink and FlowHandle::charge in each charged sub-link. FlowScope::session is only a SessionId, so the fan-out link does not keep the session alive; the SOCKS driver does, through the AppConnector it holds for as long as SocksInbound::serve runs. For a Hysteria 2 session, holding the Wire holds its quinn::Connection too.
The first flow’s Session::admit calls Ledger::bind after releasing the registry’s lock and stores the returned account in the session’s OnceLock<Arc<Account>>. Session’s Drop runs three steps in order: it cancels the session’s token, calls Account::close if an account was set, and deregisters the session from the registry. The steps of admit, the bound flag that lets later flows skip the registry’s lock, and why the account is always set before Drop can look for it are on Users, principals and sessions.
Data flow
Section titled “Data flow”From the wire to an account
Section titled “From the wire to an account”flowchart LR RT["drive: Traffic.transport_rx / tx"] --> WC1["Wire::counted"] SC["WireCounted: SOCKS stream"] --> WC1 MD["MeteredDatagram: SOCKS UDP payload"] --> WC1 QS["quinn stats: udp_rx / udp_tx"] --> CB["ConnectionBytes"] CB --> WQ["Wire::quic"] WC1 --> B["Ledger::bind at the first flow"] WQ --> B B --> ACC["Account of the UserKey"] ACC --> TK["take_usage or UsageSink"]
One byte, from add to a delta
Section titled “One byte, from add to a delta”sequenceDiagram participant D as Data plane participant W as Wire participant R as Reconcile (blocking pool) participant A as Account participant T as Take D->>W: add(up, down) R->>A: lock, then settle each live session A->>W: read() Note over A: pending += now - settled, settled = now D->>A: Session dropped: close(id) settles and removes it T->>A: lock: pending moves to reported T-->>T: one UsageDelta per user with bytes due
A byte reaches pending by exactly one of two paths, both under the account’s lock:
- Reconcile.
Ledger::reconcileruns once per sampler tick:Sampler::tickcalls it first, insidespawn_blocking, on the tick ofDEFAULT_SAMPLE_INTERVAL(1 s) unlessSupervisorBuilder::sample_intervalsets another. The push task also calls it before each batch. It moves what every live session moved since its watermark into the pending total and advances the watermark. - Close. When a bound session’s last handle drops,
DropcallsAccount::close(the order is under Who holds a wire).closesettles the session one last time and removes it fromlivein the same critical section.
A take then drains pending. So a take reports what was reconciled or closed before it: bytes a live session moved since the last tick are reported by a later take, and a session that ended is reported in full by the next one.
An account’s life in active
Section titled “An account’s life in active”stateDiagram-v2 [*] --> Active: bind (first session of the user) Active --> Active: bind, reconcile, close Active --> Kept: take finds no live session Kept --> Active: bind (a new session)
A Kept account stays in accounts with its reported total, so usage_snapshot still lists the user, but a take no longer visits it. It holds nothing pending: the take that moved it out drained it, and nothing but a bind can add to it afterwards (an_idle_reported_account_is_not_visited_until_it_binds_again).
The exactly-once mechanism
Section titled “The exactly-once mechanism”Every byte a bound session’s wire has counted by the session’s close lands in exactly one take_usage result or sink batch. The pieces that make it hold:
- One watermark per live session.
settledrecords how much of that session’s wire already reachedpending. Onlysettlemoves bytes, and it computes the difference and advances the watermark in one step. - One lock per account. Reconcile, close and take all touch an account only under its
Mutex. A reconcile and a close of the same session cannot both settle the same bytes: whichever runs second sees the advanced watermark, or no longer finds the session inlive. - Close removes the watermark with the session. Once
closehas folded a session in, the ledger no longer references its wire. The count therefore covers what the wire shows at the session’s close. - Take resets what it reports.
std::mem::takeemptiespendingand the same amount is added toreported, under the same lock. - Consistent lock order. Every path that holds both locks takes the ledger’s first and then one account’s (
bind,take,totals);reconcilereleases the ledger’s lock before it takes any account’s, andclosetakes only the account’s. No path takes an account’s lock and then the ledger’s. - Bind and take are serialised.
bindinserts the session into the account and the account intoactiveunder the ledger’s lock, which a take holds for its whole walk. A take therefore runs wholly before a bind (the session is not in the ledger yet, and its watermark will start at zero) or wholly after it (the account is inactivewith the session). - An account leaves
activeonly when it cannot owe anything.takedrops an account fromactiveonly after draining it and only whenliveis empty. A later session of the user binds again and puts it back. - A removed user keeps their account. A removal revokes principals and applies the inbound’s removal policy to the sessions; the account and its key stay. Sessions a
Keeppolicy spares keep counting under the same user, and whatever the removed user’s sessions moved is reported by the next take after it was reconciled or closed.
Taking usage
Section titled “Taking usage”Pull: take_usage
Section titled “Pull: take_usage”let mut every = tokio::time::interval(Duration::from_secs(60));loop { every.tick().await; let deltas: Vec<UsageDelta<MyUserId>> = supervisor.take_usage(); // Deliver them. On failure keep them here: take_usage will not // return these bytes again.}- At most one delta per user, none for a user with nothing due, in no particular order.
- The result lags live sessions by up to one sampler tick: what a live session moved since the last tick waits for a later take. An ended session’s bytes are all in the next take.
- Two takes racing each other, or a take racing a sink, split the pending totals between them: each byte goes to whichever drained its account first.
usage_snapshot()returns oneUsageTotalper user that ever bound a session, zero totals included, computed from the watermarks and the live wires at the time of the call. It marks nothing, so takes report the same bytes whether or not it runs. Like a take, though,totalsholds the ledger’s lock for its whole walk over every account ever bound, and reads every live session’s wire under it (a Hysteria 2 wire read takes that QUIC connection’s lock). Concurrent takes and first-flow binds wait for the walk, so poll it at the rate you need and no faster.
Live wire bytes, unlagged and unmarked, are also visible per session:
| Call | Returns | Fields | Path |
|---|---|---|---|
Supervisor::sessions() (async) |
Vec<SessionInfo<U>> |
id, inbound, source, user: Option<U>, up, down, started |
A command on the actor, which maps each UserKey to U with its own UserKeys::id, not with UsageBook::names |
Tracker::sessions() |
Vec<SessionStats> |
id, inbound, source, user: Option<UserKey>, user_label, up, down, started |
Direct, without the actor |
Both come from Sessions::stats, which lists the sessions under the registry’s lock and reads each Wire::read only after releasing it, since a QUIC wire read takes the connection’s lock and opens and admissions must not wait on that. up and down are the wire’s counts since the session opened, bound or not. The registry is described on Users, principals and sessions.
Push: usage_sink
Section titled “Push: usage_sink”SupervisorBuilder::usage_sink(sink, interval) stores a boxed closure (PushUsage<U>). SupervisorBuilder::start applies the first spec, starts the sampler, and then, if a sink was set, spawns the push task with the UsageBook and the actor’s background_stop token. Its handle joins the actor’s background list with the sampler’s.
let mut ticker = tokio::time::interval_at(tokio::time::Instant::now() + interval, interval);// A tick missed while the sink was slow is not made up in a burst.ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);loop { let stopped = tokio::select! { _ = stop.cancelled() => true, _ = ticker.tick() => false, }; // Reconciled first, off the async workers, so the batch holds // everything moved up to now rather than up to the last sampler tick. let ledger = book.ledger.clone(); if let Err(e) = tokio::task::spawn_blocking(move || ledger.reconcile()).await { std::panic::resume_unwind(e.into_panic()); } let batch = book.take(); if !batch.is_empty() { sink.report(&batch).await; } if stopped { return; }}| Step | Detail |
|---|---|
| First batch | One interval after start: interval_at does not fire at once. |
| Each tick | Reconcile on the blocking pool, then take, then report unless the batch is empty. Reconciling first is what makes a batch cover its interval exactly, rather than stopping at the last sampler tick. |
| A slow sink | The task awaits report before it waits for the next tick. With MissedTickBehavior::Delay, ticks missed meanwhile are not fired in a burst; the next batch holds more. The data plane never waits on the sink (a_slow_usage_sink_does_not_stall_the_data_plane). |
| Stop | When background_stop is cancelled, the task reconciles and takes once more, reports a non-empty batch, and returns. |
A sink that wants to hand batches to another task can forward them over a channel. The future report returns may borrow the batch, but the batch lives only until that future completes, so a sink that passes it on copies it:
use std::future::Future;use std::time::Duration;
use etemenanki_supervisor::Supervisor;use etemenanki_supervisor::entity::usage::{UsageDelta, UsageSink};use tokio::sync::mpsc;
struct Forward(mpsc::Sender<Vec<UsageDelta<MyUserId>>>);
impl UsageSink<MyUserId> for Forward { fn report(&self, batch: &[UsageDelta<MyUserId>]) -> impl Future<Output = ()> + Send { let (tx, batch) = (self.0.clone(), batch.to_vec()); async move { // Waits for room, so a slow consumer delays the next batch // instead of dropping this one. If the receiver is gone, the // batch is lost: the supervisor never offers it again. let _ = tx.send(batch).await; } }}
// `MyUserId` stands for the front end's `UserId` type, and `spec` for the// first `Spec<MyUserId>` it runs.let (tx, mut rx) = mpsc::channel(8);let (supervisor, _report) = Supervisor::builder() .usage_sink(Forward(tx), Duration::from_secs(60)) .start(spec) .await?;tokio::spawn(async move { while let Some(batch) = rx.recv().await { // Deliver `batch`, keeping it until it is delivered. }});The push task at shutdown
Section titled “The push task at shutdown”Actor::shutdown(grace) cancels background_stop only after every connection has ended, so their final bytes are in the ledger before the push task’s final reconcile:
sequenceDiagram participant C as Caller participant Act as Actor participant TT as TaskTracker participant Smp as Sampler task participant P as Push task participant Sk as UsageSink C->>Act: shutdown(grace) Act->>TT: stop listeners, cancel probes, close, wait up to grace Act->>TT: cancel root, wait for every task Note over TT: each ending Session folds its remainder in Act->>Smp: cancel background_stop Act->>P: (the same token) Smp->>Smp: final tick: reconcile, publish P->>P: reconcile, take P->>Sk: report(last batch) unless empty Act->>Act: await both tasks Act-->>C: shutdown returns
shutdownawaits the push task, so areportfuture in progress, including the one for the final batch, completes beforeshutdownreturns. A sink that talks to a remote service should bound its own call with a timeout.- When every
Supervisorhandle is dropped without ashutdown, the actor’s command loop ends and runsshutdown(Duration::ZERO), which stops the background tasks in the same order. take_usagekeeps working aftershutdown. Ina_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown,take_usageaftershutdownis empty, because the push task’s last take drained the ledger.
| Operation | Runs | Visits | Locks held |
|---|---|---|---|
Wire::add |
Per relay step (drive), per read or write (WireCounted), per datagram (charged sub-links) |
Nothing | None: up to two relaxed atomic adds |
| Counting a QUIC session | Inside quinn |
Nothing on the supervisor’s side | None |
Ledger::bind |
Once per bound session, inside AppConnector::connect on the async worker admitting its first flow |
One account | Ledger, then that account; it waits while a take or totals holds the ledger’s lock |
Account::close |
Once per bound session, when it ends | One account | That account |
Ledger::reconcile |
Every sampler tick and before every sink batch, on the blocking pool | Every live session of every active user | Ledger only while listing; one account at a time |
Ledger::take |
Each take_usage and each sink batch |
One account per active user (a live session or bytes pending) | Ledger for the walk, one account at a time |
Ledger::totals |
Each usage_snapshot |
Every account ever bound, and each live session’s wire | Ledger for the walk, one account at a time |
A take’s cost follows active users, not sessions, flows or bytes; a reconcile’s follows live sessions, and it runs in the background. Because a take and totals hold the ledger’s lock for their whole walk, and bind needs that lock, their walk time is also how long a first flow’s admission can wait on an async worker.
Only take removes an account from active. A reconcile settles each account in active that has a live session; an account without one costs a lock and an empty loop.
The benchmark supervisor/benches/tracking.rs measures a take and a reconcile over a server-sized ledger (cargo bench -p etemenanki-supervisor --bench tracking). Its ledger(per_user) helper binds per_user sessions for each of USERS = 50,000 users straight on a Ledger, numbering them SessionId::new(user * per_user + s), and returns the ledger, the wires and the accounts. No session is ever closed, so every account keeps its live sessions and stays in active through every measured take.
take_usage, 50k active users: ledgers with 1, 10 and 20 live sessions per user (50,000, 500,000 and 1,000,000 sessions). Before every measured take, each session’s wire moves 1,500 bytes each way and the ledger is reconciled, the worst case. The rows are expected to stay flat, since a take visits one account per user.reconcile (sampler tick), 50k active users: the same ledgers, 1,500 bytes each way per session before each measured reconcile. Its cost grows with the sessions.
Both groups use sample_size(20) and BatchSize::PerIteration. The file’s third group, the per-byte cost of a metered flow, belongs to Tracking.
Memory: an Account holds an Arc, a mutex, a map of live sessions and two pairs of u64; a LiveSession holds an Arc<Wire> and a pair of u64, and Account::close removes it when its session ends.
Usage and StatsSnapshot
Section titled “Usage and StatsSnapshot”The sampler publishes payload statistics once per tick, and per-user usage is kept separately. They measure different things, and neither can be derived from the other:
Usage (take_usage, UsageSink, usage_snapshot) |
StatsSnapshot (Tracker::stats) |
|
|---|---|---|
| Bytes | Wire bytes of the client’s side of a session, protocol overhead included | Payload each flow moved to and from its outbound (Metered, MeteredDatagram) |
| Unit counted | A session, bound to one user at its first flow | A flow: a stream, or one sub-link of a UDP association |
| Keyed by | The front end’s U |
UserKey, with the user’s label (UserStats) |
| Anonymous traffic | Not counted anywhere | Counted in total and per inbound and outbound, but not per user |
| Form | Deltas drained by a take, or running totals | Totals since start plus rates over the last tick, replaced every tick on a watch channel |
| Which users | Every user with bytes due (take) or ever bound (snapshot) | Users with a live flow or traffic in the last tick |
| Exactly once | Yes, between takes and a sink | Not a delta: each tick republishes totals |
| Served by the REST API | No | Yes, without the per-user rows |
The tests show the gap. In mux_sub_flows_and_their_session_count_exactly_what_moved, three VLESS mux sub-flows each count exactly their payload, while the session and the user’s single delta count every byte the socket carried (mux.sent and mux.received), mux and VLESS headers included. In a_hysteria2_session_bills_its_quic_connections_bytes, the flow counts 20,000 bytes each way and the usage more. A SOCKS UDP association is the case where the two agree on the datagrams: its sub-links’ payload is exactly what is charged to the session.
Speed limits follow the payload side: a user’s pacer is charged the payload of each of their flows (FlowHandle::add_up and add_down), not wire bytes. A user is billed for wire bytes and paced on payload.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
Every byte a bound session moves is taken exactly once (in one take_usage result or one sink batch), under its user, across any interleaving of opens, transfers, closes, removals, reconciles and takes |
Watermarks and pending under the account’s lock |
every_byte_is_taken_exactly_once |
| Nothing an unbound session moves is taken | bind only from the first flow with a user key |
every_byte_is_taken_exactly_once (10% of its sessions stay unbound) |
Once every session has ended and been taken, nothing is left, and totals equals what moved |
close folds the remainder; take drains |
every_byte_is_taken_exactly_once |
| A take never reports an empty delta | take skips a (0, 0) due |
The sum helper of tests/unit/usage.rs asserts it on every take |
| Takes and reconciles racing sessions on other threads still count every byte once | Lock order; close and reconcile both settle under the account’s lock | takes_racing_live_sessions_report_every_byte_once |
A take reports a live session’s bytes only once reconciled; totals shows them before that; a close folds the rest |
take drains pending only |
a_take_reports_live_bytes_once_they_are_reconciled |
An account with no live session and nothing pending leaves active, and returns on the next bind |
take’s retain |
an_idle_reported_account_is_not_visited_until_it_binds_again |
| Given the same key after a removal, a user continues in the same account | Ledger::bind finds it in LedgerInner::accounts, which keeps every account |
every_byte_is_taken_exactly_once (its Remove step re-adds the user under the same key with a new principal) |
| A user removed and added back gets the same key | build/users.rs → UserKeys::key keeps each id’s key |
Not by the proptest, which builds its principals with fixed keys (UserKey::new(u + 1)) and never uses UserKeys |
| A session binds to the first principal it admits | Session::admit’s bound flag and the registry lock |
the_first_flow_binds_the_session_to_its_principal |
| A session whose first flow presents a revoked principal is refused and stays unbound | Session::admit checks revoked() under the registry’s lock |
a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy pins the PermissionDenied refusal, the cancelled token and user == None in stats. The test does not inspect the ledger: that nothing reaches it follows from admit returning before Ledger::bind |
| A listener whose token is cancelled opens no session | Sessions::open checks stop under the registry’s lock |
a_stopped_listener_opens_no_session |
| A session’s wire is readable while it lives and gone with its last handle | Sessions::stats, Session’s Drop |
a_session_is_listed_with_its_bytes_until_its_last_handle_drops |
| A stream session’s usage is exactly the bytes its socket carried after the transport | drive feeds transport_rx and transport_tx |
mux_sub_flows_and_their_session_count_exactly_what_moved |
| SOCKS UDP payload is charged to the association’s session, on top of the control stream | charging_datagrams, FlowScope::charge |
a_udp_association_counts_each_sub_link_and_charges_its_session |
A Hysteria 2 session bills its QUIC connection’s UDP bytes; usage_snapshot and take_usage agree |
Wire::quic |
a_hysteria2_session_bills_its_quic_connections_bytes |
| A Hysteria 2 user admitted by password is billed under that user | Principals from the user set | a_hysteria2_user_set_admits_by_password |
With a 200 ms interval and steady traffic, batches arrive 0.9 to 2 intervals apart, and together with the one after shutdown they sum to exactly what the user moved |
push_usage, shutdown order |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
| A slow sink never slows a connection, and loses no byte | The push task is its own; Delay ticks |
a_slow_usage_sink_does_not_stall_the_data_plane |
Failure paths and cancellation
Section titled “Failure paths and cancellation”- A panic in a reconcile. The sampler and the push task both run
reconcileinspawn_blockingand re-raise a panic from it withstd::panic::resume_unwind, which ends that task. If the sampler’s task ends this way, nothing reconciles live sessions on a tick any more: in a front end without a sink, a take then sees a live session’s bytes only after the session closes. - A panic in
report. It ends the push task; the batch it held was already taken and is not offered again. - An unpublished key.
UsageBook::namepanics with<key> was bound before it was published.commit_keysruns before anything that could let a session bind a new key, so this marks a bug in the supervisor, not a runtime condition. - Cancellation of a session. Closing a session (a policy,
Supervisor::close, a removed inbound, shutdown) cancels its token; its tasks end and drop their handles, and the last handle’sDropfolds the remainder in. No path ends a bound session withoutAccount::close, because the fold is inDrop. A session cancelled before its first flow never binds:admitrefuses every flow withthe session was closed. - A first flow whose dial fails. The session was bound before the dial, in
AppConnector::connect, so it stays bound and its bytes count toward its user; the failure reaches the protocol core as a failed dial. - A step that fails. The bytes of the runtime step that ended in the error are not added to the wire (see Where the data plane feeds a wire); everything the earlier steps counted stays.
- A listener stopped mid-accept.
Sessions::openrefuses to open a session once the listener’s token is cancelled (a_stopped_listener_opens_no_session), and a Hysteria 2 connection handshaken at that moment gets an already-cancelled token and a connector with no session. Neither reaches the ledger. - Shutdown. The push task’s final reconcile and take run after every connection has ended; see The push task at shutdown. The ledger itself is not persisted: once the supervisor is gone, a new one starts from zero.
Limits
Section titled “Limits”| Item | Value | Defined in | Effect on usage |
|---|---|---|---|
DEFAULT_SAMPLE_INTERVAL |
1 s, unless SupervisorBuilder::sample_interval sets another |
track/sampler.rs |
How long a live session’s bytes may wait before a take sees them |
| First sink batch | One interval after start |
push_usage (interval_at) |
Nothing is pushed before the first interval has passed |
| Missed push ticks | MissedTickBehavior::Delay |
push_usage |
A late batch covers the longer interval; ticks are not made up in a burst |
| Batch size | At most one delta per active user | Ledger::take |
Grows with users, not sessions |
| Accounts | One per user that ever bound a session | LedgerInner::accounts |
Kept for the supervisor’s life |
| Published names | One per UserKey ever handed out |
UsageBook::names |
Kept for the supervisor’s life |
| Counters | u64 per direction per session, and per account |
Wire, AccountState |
Counts since the session opened (a wire) or since the supervisor started (an account’s reported), each up to u64::MAX bytes (about 18.4 EB) per direction |
Unit tests in supervisor/tests/unit/usage.rs (compiled into entity/usage.rs as its tests module):
| Test | What it drives | Behaviour it pins |
|---|---|---|
every_byte_is_taken_exactly_once |
A proptest of 1 to 199 operations over 4 users: open a session (weight 2; bound to a user 90% of the time), transfer 0 to 9,999 bytes each way on a live session (weight 6), close one (2), remove a user under UserRemovalPolicy::Close and add them back with a new principal under the same key (1), reconcile (2), take (2) |
After closing everything and a last take, the sum of all deltas equals what the bound sessions moved, per user; nothing is left to take; totals agrees |
takes_racing_live_sessions_report_every_byte_once |
8 worker threads, each opening 200 sessions and adding (3, 7) 50 times to each, spread over 3 users; one thread reconciling and one taking in loops until the workers finish |
Exact per-user sums under real concurrency |
a_take_reports_live_bytes_once_they_are_reconciled |
One session, add(10, 20), a take, a reconcile, add(1, 1), a take, a drop, a take |
The first take is empty while totals already shows (10, 20); the second reports (10, 20); the close folds (1, 1) into the third |
an_idle_reported_account_is_not_visited_until_it_binds_again |
A session that ends and is taken, then a new session of the same user | The account leaves active; it returns on the new bind; totals sums both sessions |
End-to-end tests in supervisor/tests/tracking.rs, each running real supervisors on loopback sockets:
| Test | Setup | Behaviour it pins |
|---|---|---|
mux_sub_flows_and_their_session_count_exactly_what_moved |
One VLESS mux connection with three sub-flows to a TCP echo; sampler every 50 ms | Each flow counts its payload; the session counts every byte the socket carried; the user’s one delta equals that, and a second take is empty |
a_udp_association_counts_each_sub_link_and_charges_its_session |
A SOCKS UDP association as alice sending 4, 6 and 5 bytes to two UDP echoes routed to two outbounds |
One flow per outbound; the session and the delta are (26 + 15, 14 + 15): greeting 3, authentication 13 and request 10 bytes up, method 2, status 2 and reply 10 bytes down, then 15 bytes of payload each way |
a_hysteria2_session_bills_its_quic_connections_bytes |
A Hysteria 2 server admitting alice by account, dialled by a second supervisor’s Hysteria 2 outbound; a 20,000-byte echo |
The flow counts exactly 20,000 each way; the session and the usage count more; usage_snapshot and take_usage agree |
a_hysteria2_user_set_admits_by_password |
A Hysteria 2 inbound with Hysteria2UserAuth::Password, a client with the right password and one with another |
The admitted user is billed under their id; the other client is refused |
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown |
A SOCKS inbound with a recording sink every 200 ms; a 4-byte echo every 20 ms for 5.5 intervals | At least 4 batches, spaced between 0.9 and 2 intervals apart; after shutdown(Duration::ZERO) the batches sum to (26 + echoed, 14 + echoed); usage_snapshot shows the same total; take_usage is empty |
a_slow_usage_sink_does_not_stall_the_data_plane |
The same inbound, a sink that sleeps 1 s per report, an interval of 50 ms, 11-byte echoes for 1.5 s | Every echo round trip stays under 250 ms; the sink sees at most 2 batches meanwhile; the final sum is exact |
Both sink tests use the file’s Recorder sink and its socks_with_sink(sink, interval) helper, which starts a SOCKS inbound admitting alice from an account set. Recorder::report pushes (Instant::now(), batch.to_vec()) onto a shared Vec when it is called, before its future is polled, and returns tokio::time::sleep(delay) as that future: Duration::ZERO in the first test, 1 s in the second. Recorder::sum adds up every recorded delta.
Session tests in supervisor/tests/unit/session.rs that bear on usage: a_session_is_listed_with_its_bytes_until_its_last_handle_drops, the_first_flow_binds_the_session_to_its_principal, a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy and a_stopped_listener_opens_no_session (a cancelled listener token makes Sessions::open return None and leaves stats empty). The rest of that file is described on Users, principals and sessions.
No test covers UsageBook::name’s panic, a panicking sink, a sink installed on a supervisor dropped without shutdown, or UserKeys handing a re-added user the same key. Add one when you touch those paths. For how to run the suites, see Testing.