Skip to content

Metering and rate limiting

Source files: 13 · checked against katana v3.0.1 · Etemenanki 596916d
  • katana/src/meter.rs
  • katana/src/traffic.rs
  • katana/src/connector.rs
  • katana/src/outbound/proxy.rs
  • katana/src/serve.rs
  • katana/src/manager/node.rs
  • katana/src/manager/mod.rs
  • katana/src/api/mod.rs
  • katana/tests/unit/meter.rs
  • katana/tests/unit/traffic.rs
  • katana/tests/unit/connector.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/wake.rs

katana counts and paces traffic in one place: the outbound stream it hands to the kernel’s per-connection runtime. src/meter.rs wraps that stream in Metered, which bills every byte to the user’s UserCounter and holds the flow back while the user’s TokenBucket is in debt. The same Gate sits inside the UDP fan-out, so datagrams are paced by the same bucket.

This page is for contributors who change the speed limit, the byte counters or the outbound path. It covers the three types involved, why billing happens on the outbound side, why the long-run rate does not depend on how the runtime chunks its reads and writes, and the tests that pin each of these. How the counters are reported to the panel is on Traffic accounting; where a user’s rate comes from, from an operator’s point of view, is on Speed limits and connection limits.

Component Where Does Leaves to others
TokenBucket src/traffic.rs Keeps one byte balance per user: refills at rate, caps at one second of rate, goes negative on a charge larger than the balance, and answers “when is the debt repaid”. Waiting. It never sleeps outside tests.
UserCounter src/traffic.rs Holds the up and down byte totals and the user’s bucket. One per registered user, shared by every flow of that user through an Arc. Reporting and reconciling the totals (NodeTraffic).
Gate src/meter.rs Answers “may this flow move bytes now”: not retired, and the bucket not in debt. Bills a transfer after it moved. Deciding who is retired (Admission).
Metered<S> src/meter.rs Wraps an outbound byte stream: writes are upload, reads are download, each gated before and billed after. Moving the bytes (the inner stream) and scheduling (the runtime).
KatanaConnector::connect src/connector.rs Builds one Gate per flow from the admitted user’s counter and lease, and wraps the stream or the UDP fan-out with it. Routing and audit decisions.

katana has no relay loop of its own. The Etemenanki runtime moves every byte of a connection in one task, and what katana controls is each outbound the runtime opens through the connector. So the outbound is the only place where katana can see bytes cross, and it is where metering lives.

src/traffic.rs
pub struct TokenBucket {
rate: u64,
state: Mutex<BucketState>,
}
struct BucketState {
tokens: f64,
last: Instant,
}
impl TokenBucket {
pub fn new(rate: u64) -> Self;
pub fn charge(&self, n: usize) -> Option<Instant>;
pub fn ready_at(&self) -> Option<Instant>;
#[cfg(test)]
pub async fn consume(&self, n: usize);
fn refill(&self, st: &mut BucketState) -> Instant;
fn repaid_at(&self, st: &BucketState, now: Instant) -> Option<Instant>;
}

Mutex is parking_lot::Mutex and Instant is tokio::time::Instant, which is what lets the tests run on paused time.

Item Meaning
rate Bytes per second. 0 means unlimited: charge and ready_at return None at once, without taking the lock. Fixed at construction.
tokens The balance in bytes, as f64. Negative while the bucket is in debt. Starts at rate, so a new bucket is full.
last When the balance was last brought up to date.
refill Credits elapsed × rate since last and clamps the balance to at most rate. That clamp is the burst: one second of rate, however long the bucket sat idle.
charge(n) Refills, subtracts n in full, and returns repaid_at. The subtraction happens even when n exceeds the balance.
ready_at() Refills and returns repaid_at without charging.
repaid_at None when tokens >= 0; otherwise now + (-tokens / rate) seconds, the instant the refill brings the balance back to zero.
consume(n) Test-only: charge(n), then sleep_until the returned instant.

Every method that reads the balance calls refill first under the same lock, so the balance is always current when it is compared or changed.

src/traffic.rs
pub struct UserCounter {
pub uid: i64,
up: AtomicU64,
down: AtomicU64,
pub rate: u64,
pub bucket: TokenBucket,
}
impl UserCounter {
pub(crate) fn new(uid: i64, rate: u64) -> Self;
pub fn add_up(&self, n: u64);
pub fn add_down(&self, n: u64);
pub fn up(&self) -> u64;
pub fn down(&self) -> u64;
pub fn commit_reported(&self, up: u64, down: u64);
}

add_up and add_down are fetch_add with Ordering::Relaxed. The counter’s rate is the user’s effective rate in bytes per second. build_user_entries (src/manager/mod.rs) computes it for every user it registers (a user whose UUID does not parse is skipped) with

src/traffic.rs
pub fn determine_rate(node_bps: u64, user_bps: u64) -> u64;

which returns the smaller of the two non-zero limits, and 0 (unlimited) when both are 0.

src/meter.rs
pub struct Gate {
counter: Arc<UserCounter>,
wait: Option<Pin<Box<Sleep>>>,
retired: Pin<Box<WaitForCancellationFutureOwned>>,
}
impl Gate {
pub fn new(counter: Arc<UserCounter>, retired: CancellationToken) -> Self;
pub fn poll_open(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
pub fn sent(&self, n: usize);
pub fn received(&self, n: usize);
}
Field or method Meaning
counter The user’s counter, shared with every other flow of that user.
wait The one tokio::time::Sleep this gate parks on while the bucket is in debt. Both directions of the flow share it.
retired CancellationToken::cancelled_owned() on the user’s lease, boxed and pinned once in new.
poll_open(cx) Ready(Ok(())) when the flow may move bytes; Ready(Err(_)) once the user is retired; Pending while the bucket is in debt, with cx registered on the retirement and on the timer. Every answer, Ready(Ok(())) included, leaves cx registered on the retirement.
sent(n) counter.add_up(n), then counter.bucket.charge(n).
received(n) counter.add_down(n), then counter.bucket.charge(n).

Because poll_open polls the retirement even when it lets the transfer through, a flow that then parks on its inner stream, such as a read waiting on a silent destination, is still woken when the user is retired. The next poll returns the ConnectionAborted error.

sent and received discard the instant charge returns. The debt a transfer leaves is paid by whatever calls poll_open next, on this flow or any other flow of the same user.

Gate is not Clone and holds no lock. Its two pinned futures are boxed so that Gate, and with it Metered<S> for any S: Unpin, stays Unpin; Metered relies on that when it calls self.get_mut().

src/meter.rs
pub struct Metered<S> {
inner: S,
gate: Gate,
}
impl<S> Metered<S> {
pub fn new(inner: S, gate: Gate) -> Self;
}
impl<S: AsyncRead + Unpin> AsyncRead for Metered<S> {
fn poll_read(
self: Pin<&mut Self>,
cx: &mut Context<'_>,
buf: &mut ReadBuf<'_>,
) -> Poll<io::Result<()>>;
}
impl<S: AsyncWrite + Unpin> AsyncWrite for Metered<S> {
fn poll_write(
self: Pin<&mut Self>,
cx: &mut Context<'_>,
data: &[u8],
) -> Poll<io::Result<usize>>;
fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
}

The connector instantiates it as Metered<OutboundStream>, the associated Stream type of impl Connector<Flow<UserTag>> for KatanaConnector. OutboundStream (src/outbound/proxy.rs) is a plain TCP stream, a boxed proxy-client stream, or a WireGuard stream.

Method Gated Billed
poll_read Yes, before the inner read received(n) for the n bytes the inner read added to buf, if n > 0
poll_write Yes, before the inner write sent(n) for the n bytes the inner write accepted, if n > 0
poll_flush No No
poll_shutdown No No

Metered implements only these four methods. tokio’s default poll_write_vectored forwards the first non-empty slice to poll_write, so a vectored write is gated and billed like any other.

The runtime decodes the inbound protocol, then forwards payload to the outbound and payload from it back to the client. Metered wraps the outbound, so it sees the bytes after the inbound protocol took its framing and encryption off, and before the outbound protocol puts its own on.

flowchart LR
  client["client"] --> transport["inbound transport: TCP, TLS, WS, gRPC"]
  transport --> core["protocol core: VMess, VLESS, Trojan, SS"]
  core --> runtime["ProxyServerRuntime"]
  runtime --> metered["Metered: Gate, UserCounter"]
  metered --> outbound["OutboundStream: TCP, proxy client, WireGuard"]
  outbound --> target["destination"]

That position is chosen for three reasons:

  • The payload is plain there. A user is billed for what they asked to move, not for TLS records, WebSocket frames, gRPC framing or VMess chunk headers. Two users with the same traffic pay the same whatever inbound they connect through.
  • It is the only hook. The runtime owns the inbound side; katana supplies only the connector and what it returns.
  • Every flow passes through it. Each TCP request, each mux sub-flow, each Hysteria stream and each UDP association reaches KatanaConnector::connect, which is where the gate is built.

The diagram shows the stream listeners (serve_stream in src/serve.rs). A Hysteria 2 listener takes a different inbound path: the Etemenanki Hysteria inbound runs one runtime per proxy stream and one per connection’s UDP, and each of them opens its outbounds through a KatanaConnector built by run_hysteria. The outbound side, and so the metering, is the same.

Metered bills the plaintext it is given. If the outbound is itself a proxy client (VMess, VLESS, Shadowsocks and so on), the client’s own framing is added inside inner and is not billed.

The gate is built in the connector, the one place every flow of every inbound passes:

src/connector.rs
impl Connector<Flow<UserTag>> for KatanaConnector {
type Stream = Metered<OutboundStream>;
type Datagram = FanOut;
type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;
}
impl Admission {
pub fn admit(&self, tag: &UserTag) -> Option<(Arc<UserCounter>, CancellationToken)>;
}

KatanaConnector::connect admits the flow’s user with Admission::admit, which returns the user’s Arc<UserCounter> and a CancellationToken lease, and then builds one gate per flow:

src/connector.rs (excerpt)
let gate = Gate::new(counter, lease);
  • A UDP flow gets FanOut::new(self.disp.clone(), gate, tag.uid, source).
  • A TCP flow is routed and audited first. A flow sent to Outbound::Block or forbidden by an audit rule is refused before anything is dialled, so it is never billed. Otherwise the dialled stream is wrapped: Metered::new(dial.await?, gate).

Each flow has its own Gate and its own wait timer. All gates of one user hold the same Arc<UserCounter>, because NodeTraffic::lookup returns the one counter registered for that credential and uid. That is what makes the bucket per user rather than per flow. The registry belongs to one node (each NodeManager owns its own NodeTraffic), so a user served by two nodes of the same katana process has one bucket on each.

The runtime writes a forward range to the outbound with poll_write and, when the write accepted bytes, calls poll_flush in the same step. Each outbound is polled with its own KeyWaker (concepts/src/wake.rs), so a wake-up from the gate’s timer reschedules exactly that outbound.

sequenceDiagram
  participant R as runtime task
  participant M as Metered
  participant G as Gate
  participant C as UserCounter
  participant B as TokenBucket
  participant S as inner stream
  R->>M: poll_write(cx, data)
  M->>G: poll_open(cx)
  G->>G: poll retired (registers cx)
  G->>B: ready_at()
  B-->>G: Some(at), bucket in debt
  G->>G: wait = sleep_until(at), poll it
  G-->>M: Pending
  M-->>R: Pending, nothing moved
  Note over R,G: the timer fires and wakes this outbound's KeyWaker
  R->>M: poll_write(cx, data)
  M->>G: poll_open(cx)
  G->>G: poll retired again
  G->>G: timer fired, clear wait
  G->>B: ready_at()
  B-->>G: None, out of debt
  G-->>M: Ready(Ok)
  M->>S: poll_write(cx, data)
  S-->>M: Ready(Ok(n))
  M->>G: sent(n)
  G->>C: add_up(n)
  G->>B: charge(n), result ignored
  M-->>R: Ready(Ok(n))
  R->>M: poll_flush(cx), not gated

A read is the mirror image: poll_open, then the inner poll_read, then received(n) for the bytes that landed in the ReadBuf. A read that returns end of stream (n == 0) bills nothing.

poll_open checks the retirement first on every call, then loops over the bucket:

stateDiagram-v2
  [*] --> CheckRetired
  CheckRetired --> Aborted: lease cancelled
  CheckRetired --> AskBucket: not cancelled
  AskBucket --> Open: ready_at is None
  AskBucket --> Waiting: ready_at is Some(at)
  Waiting --> Parked: timer not due
  Waiting --> AskBucket: timer fired, wait cleared
  Parked --> CheckRetired: woken by timer or cancellation
  Open --> [*]
  Aborted --> [*]
  • Aborted returns Err(io::Error::new(io::ErrorKind::ConnectionAborted, "the user was retired")).
  • Parked is Pending with cx registered on both the cancellation future and the Sleep, so whichever happens first wakes the flow.
  • After a timer fires, the gate asks the bucket again instead of assuming it is clear. Another flow of the same user may have charged more debt meanwhile; if so, the gate arms a new timer for the new instant.

FanOut (see Connector and UDP fan-out) uses the same Gate with the same order: gate, then transfer, then bill.

Step poll_send_to poll_recv_from
Before the gate The packet is routed and audited. A blocked or forbidden packet returns Ready(Ok(buf.len())) and is dropped unbilled. Nothing.
Gate self.gate.poll_open(cx) self.gate.poll_open(cx)
Transfer The sub-link for the chosen outbound sends the packet; a sub-link still opening returns Pending. If opening the sub-link fails, the packet is dropped, reported as sent, and not billed. The sub-links are polled round-robin from next.
Bill self.gate.sent(n) on Ready(Ok(n)) self.gate.received(len) for the payload bytes that landed in the buffer

So a user’s TCP streams, mux sub-flows and UDP packets all draw from the same bucket, in both directions.

Invariant Mechanism Pinned by
One bucket per user, shared by all their flows and both directions, over TCP and UDP. Every Gate of a user holds the Arc<UserCounter> from NodeTraffic::lookup; sent and received both call counter.bucket.charge; FanOut and Metered both bill through Gate. the_limit_is_shared_by_both_directions (tests/unit/meter.rs), debt_holds_back_the_next_charge_too (tests/unit/traffic.rs), udp_is_billed_after_routing_and_blocked_packets_are_free (tests/unit/connector.rs)
Writes bill upload and reads bill download, exactly the bytes that moved. poll_write bills the n the inner write returned; poll_read bills the growth of buf.filled(). each_direction_is_billed_to_the_user (tests/unit/meter.rs), an_admitted_stream_is_billed_to_its_user (tests/unit/connector.rs)
The long-run rate does not depend on chunk size. charge takes the whole transfer even past zero; poll_open refuses to start any transfer while tokens < 0. the_limit_holds_however_the_writes_are_sized (tests/unit/meter.rs), a_chunk_larger_than_the_burst_is_still_limited (tests/unit/traffic.rs)
An idle user banks at most one second of rate. refill clamps tokens to rate. an_idle_bucket_banks_one_second_and_no_more (tests/unit/traffic.rs)
An unlimited user is never delayed. rate == 0 short-circuits charge and ready_at to None. token_bucket_unlimited_is_instant (tests/unit/traffic.rs)
No Pending after bytes moved. The gate is consulted before the inner call; after the inner call returns Ready, Metered only bills and returns Ready. A read never loses data it put in buf, and a write never reports “nothing written” for bytes the inner stream took. By construction; each_direction_is_billed_to_the_user (tests/unit/meter.rs) and an_admitted_stream_is_billed_to_its_user (tests/unit/connector.rs) check that the bytes arrive intact and are billed once.
A Pending from the gate means nothing moved. ready!(this.gate.poll_open(cx))? returns before the inner stream is touched. By construction.
Retirement ends a flow even while it is parked. poll_open polls retired before anything else, registering cx on the cancellation; the retirement stays observed on every later poll because the owned cancellation future checks is_cancelled first. retiring_the_user_wakes_a_parked_read, a_retired_users_stream_refuses_to_move (tests/unit/meter.rs)
Blocked and forbidden traffic is free. TCP: connect refuses before building Metered. UDP: poll_send_to drops before the gate. udp_is_billed_after_routing_and_blocked_packets_are_free, a_forbidden_udp_destination_is_dropped_and_recorded (tests/unit/connector.rs)
A retired user is not billed for flows they cannot open. Admission::admit returns None for a departed or rebound credential, and connect returns a refusal before a gate exists. a_user_the_registry_does_not_know_is_refused, a_credential_rebound_to_another_uid_is_refused (tests/unit/connector.rs)

Why debt makes the rate independent of chunk size

Section titled “Why debt makes the rate independent of chunk size”

The runtime does not read or write in fixed sizes. It reads from the outbound into scratch[..max] with max = (staging.room() - Core::STAGING_RESERVE).min(BUF_SIZE), and it writes whatever forward range the protocol core produced. BUF_SIZE differs per protocol core. A limiter whose behaviour depends on those sizes would give users different speeds on different protocols.

Consider the alternative, “wait until the bucket holds n tokens, then move n”. With a burst cap of one second, a transfer larger than rate can never find enough tokens, so such a limiter must let it through after at most one refill period, and a flow that moves large chunks runs faster than rate.

TokenBucket avoids that by letting the balance go negative:

  1. A transfer may start whenever tokens >= 0.
  2. After it moves, its full size is subtracted, even past zero.
  3. No transfer of that user starts again until the refill brings tokens back to zero, at rate bytes per second.

Every byte is therefore paid for at rate, whether it moved in one piece or in many. The only slack is the burst (at most rate bytes banked) plus the transfers that had already passed the gate when the balance crossed zero; their cost is paid by the transfers after them.

Worked example, from the_limit_holds_however_the_writes_are_sized, with rate = 10_000 and a full bucket:

Step Balance before Transfer Balance after Next transfer may start
Write 30,000 bytes 10,000 moves at once −20,000 2 s later
Write 1 byte 0 (after 2 s) moves −1 0.0001 s later

Gate keeps a single wait because both directions of one Metered are polled from the same runtime task with the same per-outbound waker (slot.waker in concepts/src/runtime.rs). A Sleep stores one waker; since read and write register the same one, neither direction can be left without a wake-up. When the timer fires, the KeyWaker queues that outbound and wakes the runtime task. On that poll the runtime retries the write still waiting at the front of its effect queue, and reads the outbound because its key was queued, so both directions see the same bucket state.

Situation What Metered or FanOut returns What happens next
The user’s lease is cancelled (they left the user list, their credential moved to another uid, or Admission::retire_all ran because the listener is going away) Err with io::ErrorKind::ConnectionAborted and the message the user was retired, from every later read, write, send or receive The runtime fails that outbound. On a stream listener the connection itself also ends through the lease slot; see Admission and user tables. A Hysteria 2 connector has no lease slot, so there the flows end only through their outbounds refusing to move.
The inner stream fails The inner error, unchanged Nothing is billed for that call.
The inner read returns end of stream Ready(Ok(())) with nothing added Nothing is billed.
The bucket is in debt Pending The flow resumes when the timer fires or the lease is cancelled, whichever comes first.
poll_flush or poll_shutdown while in debt or retired Forwarded to the inner stream Flushing bytes that were already accepted, or closing the stream, never waits on the bucket and is not refused after retirement.

Dropping a Metered drops its Gate, which drops the pending Sleep and the cancellation future. Neither holds a task or a lock, so cancellation needs no clean-up. The bucket’s Mutex is held only for the arithmetic inside charge and ready_at, never across an .await or a poll.

A new rate reaches the node through a panel poll: the regular one, or the one NodeManager::apply_static (src/manager/node.rs) runs at once when a config edit changes [node.api], for example the speed_limit override. The panel client keeps its own copy of that override, so such an edit builds a new PanelClient before the poll. The new client holds no ETags, so the panel answers that poll with the full user list rather than 304 Not Modified, and every changed rate shows up in it.

A rate change is not a retirement. NodeTraffic::prepare gives the user a fresh UserCounter, with a fresh bucket that starts full, because a bucket’s rate is fixed at construction. Existing gates keep the Arc to the old counter, so:

  • flows opened before the change keep billing the old counter and keep the old rate until they end;
  • flows opened after it use the new counter and the new rate;
  • for a while the user draws from two buckets at once.

The old counter drains and is still reported; see Traffic accounting.

Name Value Where Meaning
Burst rate bytes (one second) TokenBucket::refill The most a bucket can bank. A new bucket starts with exactly this.
Unlimited rate == 0 TokenBucket::charge, ready_at No pacing; counters still count.
MBPS_TO_BPS 1_000_000.0 / 8.0 src/api/mod.rs Panels give limits in Mbps; mbps_to_bps multiplies by this, so 1 Mbps is a rate of 125,000 bytes per second. A value of 0 or below becomes 0, unlimited.
MAX_SUBS 64 src/connector.rs Sub-links one UDP association keeps; they all share the association’s single Gate.

The bucket works in bytes per second. The panel clients convert every limit with mbps_to_bps before determine_rate combines the node and user limits. Which panel field feeds which limit, and the node-level override, are described in Speed limits and connection limits.

tests/unit/meter.rs, run against tokio::io::duplex pairs:

Test Pins
each_direction_is_billed_to_the_user A 5-byte write bills up = 5, a 7-byte read bills down = 7.
the_limit_holds_however_the_writes_are_sized On paused time at 10,000 B/s, one 30,000-byte write followed by a 1-byte write takes at least 2 s: the oversized write left 20,000 bytes of debt.
the_limit_is_shared_by_both_directions A 20,000-byte upload at 10,000 B/s delays the next 5-byte read by at least 1 s.
a_retired_users_stream_refuses_to_move A write after the lease is cancelled fails with ConnectionAborted.
retiring_the_user_wakes_a_parked_read A read parked on a silent peer is woken by the cancellation within 1 s and fails with ConnectionAborted.

tests/unit/traffic.rs, on the bucket alone:

Test Pins
token_bucket_unlimited_is_instant rate = 0 never waits, even for 1,000,000 bytes.
token_bucket_rate_limits After the burst is drained, the next 10,000 bytes at 10,000 B/s wait about 1 s.
a_chunk_larger_than_the_burst_is_still_limited consume(30_000) at 10,000 B/s takes exactly 2 s on paused time.
debt_holds_back_the_next_charge_too charge(25_000) returns an instant 1.5 s away, and the next 1-byte charge waits it out.
an_idle_bucket_banks_one_second_and_no_more After 60 idle seconds, charge(20_000) still leaves 1 s of debt.
determine_rate_min_nonzero The effective rate is the smaller non-zero limit.

tests/unit/connector.rs, through KatanaConnector::connect against local echo servers:

Test Pins
an_admitted_stream_is_billed_to_its_user A TCP round trip of 5 bytes bills (5, 5) to the registered counter.
udp_is_billed_after_routing_and_blocked_packets_are_free A UDP echo bills (4, 4); a packet to a blocked address is accepted and not billed.
a_forbidden_udp_destination_is_dropped_and_recorded A packet an audit rule forbids is not billed and is recorded as a hit.

The tests that measure pacing run on #[tokio::test(start_paused = true)], so the durations they assert are exact and do not depend on machine load. Two tests in tests/unit/traffic.rs use real time instead: token_bucket_rate_limits asserts a lower bound of 0.8 s, and token_bucket_unlimited_is_instant an upper bound of 50 ms. retiring_the_user_wakes_a_parked_read also runs on real time, because it only asserts that the wake-up arrives within its 1 s timeout. When you add a timing test for the bucket or the gate, use paused time.