跳转到内容

按用户的用量计费

源码文件:29 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
  • Etemenanki/supervisor/src/entity/usage.rs
  • Etemenanki/supervisor/src/entity/session.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/entity/user.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/build/users.rs
  • Etemenanki/supervisor/src/build/inbound.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/topology/spec_plan/inbound.rs
  • Etemenanki/supervisor/src/topology/outbound/udp_fanout.rs
  • Etemenanki/supervisor/src/track/mod.rs
  • Etemenanki/supervisor/src/track/metered.rs
  • Etemenanki/supervisor/src/track/sampler.rs
  • Etemenanki/supervisor/benches/tracking.rs
  • Etemenanki/supervisor/Cargo.toml
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/obfs.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/ffi/src/proxy.rs
  • Etemenanki/webclient/src/lib.rs
  • Etemenanki/webclient/src/tracking.rs
  • Etemenanki/supervisor/tests/unit/usage.rs
  • Etemenanki/supervisor/tests/unit/session.rs
  • Etemenanki/supervisor/tests/tracking.rs
  • katana/src/manager/node.rs

supervisor(监管器)为每个用户统计其客户端连接在服务端线路上产生的字节,并以增量的形式交给前端程序:由前端程序轮询(Supervisor::take_usage),或者按固定间隔推送给一个 UsageSink。已绑定会话在关闭之前统计到的每个字节,都恰好落入一个增量,不论是两种方式中的哪一种取走了它。计数在每个会话的 Wire 中完成;记账在每个 supervisor 唯一的一个 Ledger 中完成,它为每个用户维护一个账户,并为每个存活会话维护一个水位(watermark)。

本页面向修改 supervisor/src/entity/usage.rs 的贡献者,以及基于 supervisor 做计费的前端程序作者。内容包括:统计什么、不统计什么,Wire,Ledger 及其账户,恰好一次(exactly-once)的机制,拉取与推送两条路径(包括关停时的推送任务),开销,以及用量与 StatsSnapshot 中载荷数据的区别。会话本身(如何打开、如何绑定 principal(身份主体)、如何结束)见用户、principal 与会话。流、采样器和限速见追踪。

用量代码负责:

  • 为每个会话提供一个 Wire,从会话的第一个字节起统计其客户端一侧承载的字节;
  • 会话的第一条流把它绑定到某个用户时,在 Ledger 中打开该用户的账户;会话结束时,把会话的剩余字节并入这个账户;
  • 在每个采样器 tick 把每个存活会话的新字节移入其用户的待取总量,这一步不在异步 worker 上执行;
  • 以 UsageDelta 值的形式交出待取总量,每次取出(take)中每个用户最多一个,且同一个字节从不交出两次;
  • 安装了 UsageSink 时,运行一个推送任务:在其间隔的每个 tick 调和(reconcile)、取出并上报,并在关停流程等到所有连接结束之后再执行一次;
  • 报告每个用户自 supervisor 启动以来的累计总量(usage_snapshot),但不取走任何字节。

它不负责:

  • 统计载荷。每条流搬运的载荷归采样器管,最终进入 StatsSnapshot(见用量与 StatsSnapshot)。
  • 把任何数据送到面板。交出去的增量由前端程序负责投递;supervisor 不会再次提供它。
  • 持久化。账本只存在于内存中,Actor::new 为每个 supervisor 新建一个空账本。
  • 限速。限速按每条流的载荷扣减,通过追踪一页描述的 pacer 完成。
前端程序 路径 位置
katana 拉取:在其上报周期中调用 Supervisor::take_usage。它用 Supervisor::start(spec) 启动每个节点的 supervisor,因此采样器以默认的 1 s 为 tick 间隔,也没有安装 sink src/manager/node.rs → report_traffic;见 katana 流量计费
etemenanki-app 两者都不用:Instance::start 用 Supervisor::builder() 构建 supervisor,不安装 sink app/src/instance.rs
etemenanki-ffi 两者都不用:它用 Supervisor::builder().socket_options(...) 构建 supervisor,不安装 sink ffi/src/proxy.rs
REST API(etemenanki-webclient) 两者都不用:它的流量视图不含按用户的数据,其文档说明按用户的用量只能通过 supervisor 的 Rust API 获取 webclient/src/lib.rs、webclient/src/tracking.rs

REST API 本身在REST API(etemenanki-webclient)一页中介绍。

用量是每个会话的线路字节(wire bytes):服务端带宽为该用户的客户端承载的字节,包括协议开销。supervisor/src/entity/usage.rs 的模块文档这样定义它;每种会话从不同的位置向自己的 Wire 计数。

会话 Wire 计数来源 统计范围
流式入站(HTTP、Trojan、VLESS、VMess、Shadowsocks、Shadowsocks 2022)经由明文 TCP、TLS 或 WebSocket 的连接 该套接字会话的 Wire::counted(),由 run_stream_inbound 在 accept 时、传输层自身握手之前打开 serve.rs → drive:运行时产出的每个 Traffic 步骤都把 transport_rx 记为上行、transport_tx 记为下行 传输层之后、协议的分帧和加密尚未解开之前的协议流:协议握手、头部、填充和密文都计入。TLS 或 WebSocket 握手、TLS 记录和 WebSocket 帧不计入:在 drive 开始处理传输层交出的 stream 之前,没有任何东西向 wire 计数
gRPC 承载连接(carrier)上的一条 stream 每条 HTTP/2 stream 都是一个独立的会话,有自己的 Wire::counted() drive,同上 这条 stream 在 gRPC 分帧之后的字节
SOCKS 连接 该套接字会话的 Wire::counted() serve.rs → WireCounted,包在客户端 stream 外面的一层包装:每次成功的 poll_read 计入上行,每次成功的 poll_write 计入下行 控制流的每个字节:问候、认证、请求和应答,以及 CONNECT 中继的字节
SOCKS UDP 关联 同一个会话的 Wire AppConnector::charging_datagrams → FlowScope::charge → 每个子 link 的 MeteredDatagram;FlowHandle::add_up 和 add_down 也会累加到被计费的 wire 与目的地之间中继的每个数据报的载荷。这些数据报从不经过控制流,所以这是它们计入会话的唯一途径
Hysteria 2 QUIC 连接 Wire::quic(ConnectionBytes),在 run_hysteria_inbound 中构建 quinn::Connection::stats():udp_rx.bytes 记为上行,udp_tx.bytes 记为下行 quinn 在该连接上收发的每个数据报的 UDP 载荷:QUIC 握手、QUIC 分帧、TLS 和 HTTP/3 凭据交换都包括在内。IP 和 UDP 头部不计入。启用混淆时,Salamander 层(protocols/src/hysteria/obfs.rs)在 quinn 之下给每个数据报添加的 SALT_LEN(8)字节盐值也不计入,因为 quinn 看不到它
Unix socket 监听器的连接 该套接字会话的 Wire::counted(),serve_socket 服务它时不带源地址(connector_for(tag, None, carrier)) drive 或 WireCounted,与 TCP 相同 该 socket 的字节;Unix 监听器没有传输层

方向以用户为准:up 是其客户端发往目的地的字节,down 是返回给客户端的字节。

因此,Hysteria 2 是唯一一种计数包含其传输层自身开销的会话:QUIC 连接既是传输层,也是会话。测试 a_hysteria2_session_bills_its_quic_connections_bytes 固定了这一点:一次 20 000 字节的回显,在流上每个方向正好是 20 000 字节,而在会话和该用户的用量中每个方向都超过 20 000 字节。

会话在其第一条流绑定它之后才计入其用户。Session::admit 把会话绑定到它准入的第一条流的 Principal;如果这个 principal 带有 UserKey,admit 就用会话的 Wire 调用 Ledger::bind。绑定时会话的水位从 (0, 0) 开始,所以 wire 在绑定之前统计的一切都会和其余字节一起上报:流式会话的协议握手,或者 Hysteria 2 连接的 QUIC 握手和凭据交换。

AppConnector::connect 在路由或拨号任何东西之前,为每条流执行这次 admit(见路由与 plane(数据平面))。因此绑定不依赖拨号:第一条流打开失败的会话照样被绑定,它的 wire 统计到的字节(包括协议握手)都计在其用户名下。

之后的流即使出示另一个 principal,也不会重新绑定会话(the_first_flow_binds_the_session_to_its_principal)。一条 mux 连接是一个会话,所以每个子流的字节都计入打开第一条流的那个用户。

情形 为什么没有字节进入账户
从不打开流的会话:传输层或协议握手失败或超时、从不认证的 Hysteria 2 客户端,或者认证后从不打开 stream 或数据报流的客户端 Ledger::bind 只由第一条被准入的流触发。在协议核心(core)进入 established 状态(请求已解析、流已打开或正在打开)之前,drive 对每个运行时步骤最多等待 HANDSHAKE_TIMEOUT(10 s,protocols/src/core/mod.rs),超时则以 inbound handshake timed out after 10s 结束连接;这样的连接没有打开任何流,所以其 wire 统计到的字节不计入任何人
在第一条流之前就被关闭的会话 一旦会话的 token 被取消或其注册表条目已不存在,Session::admit 就会在查看 principal 之前,以 the session was closed(PermissionDenied)拒绝每一条流
第一条流出示的 principal 在绑定之前已被吊销的会话 Session::admit 以 the user was removed before the session opened a flow 拒绝它,并在到达 Ledger::bind 之前返回
以匿名身份准入的流:处于开放模式或共享模式的入站 它的 principal 是 Principal::anonymous(),其 user_key() 为 None。build/inbound.rs → user_table 把它交给没有用户集的入站(InboundSpec::users 为 None,该类型定义在 topology/spec_plan/inbound.rs):未启用认证的 SOCKS、没有账户的 HTTP、只有服务端密码的 Shadowsocks、只有服务端 PSK 的 Shadowsocks 2022,以及使用 shared_password 的 Hysteria 2
gRPC 承载连接自身的分帧 gRPC 承载连接上的每条 HTTP/2 stream 都是一个独立的会话,统计的是 gRPC 分帧之后的字节;承载连接自身的分帧不计入任何人
TUN 设备 设备不准入用户,所以它的 connector(连接器)没有会话(run_tun_inbound)
supervisor 自己打开的流 它们带有 anonymous_user()(topology/flow.rs)

按用户的载荷(包括按入站和出站统计的匿名流量)见采样得到的 StatsSnapshot。

supervisor/src/entity/usage.rs
#[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,
}
  • UsageDelta 是某个用户自上次取出以来的流量。take_usage 的结果和一个 sink 批次中,每个用户最多一个,并且不包含 up 和 down 都为零的增量。
  • UsageTotal 是某个用户自 supervisor 启动以来的全部流量,无论是否已上报。
  • 两者都以前端程序自己的 U: UserId 为键,而不是内部的 UserKey。两者都 derive 了 serde 的 Serialize 和 Deserialize,只要 U 实现了它们就生效,因此前端程序可以原样发送上报;这就是 supervisor/Cargo.toml 依赖 serde 的原因。
supervisor/src/entity/user.rs
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 就是前端程序本来对其用户的称呼:面板的数字 id,或者实现了该 trait 的 UserName(一个 Email 或一个 Username)本身。其背后的模型见用户、principal 与会话。

supervisor/src/entity/usage.rs
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;
}
  • report 在推送任务上运行,从不在数据平面上运行。慢的 sink 只会推迟它自己的下一个批次,那个批次随之覆盖更长的间隔。
  • 交给 sink 的增量不会再由 take_usage 返回。
  • report 返回 ()。supervisor 无从得知某个批次没有送达,也从不再次交出同一个批次:不能丢失批次的 sink 要自己保存批次,直到送达为止。
supervisor/src/entity/usage.rs
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);
}
方法 行为
counted() 两个计数器都为零,没有 QUIC 读取器。
quic(bytes) 同样的计数器,外加 bytes,即某条 QUIC 连接字节数的读取器。
add(up, down) 每个方向一次 Ordering::Relaxed 的 fetch_add,值为零的方向跳过。这就是数据平面上每字节的全部开销。
read() 到目前为止的 (up, down),从不减小。对 QUIC wire,它返回连接的 (received, sent) 加上计数器的值。没有任何东西向 Hysteria 2 会话的计数器累加,所以它的 read 就是连接的计数。

Relaxed 顺序就足够了,因为没有任何数据通过这些计数器发布:每个读取方只需要一个永不回退的值,而账本的水位由账户的锁保护,而不是由原子变量保护。

protocols/src/hysteria/server/inbound.rs
#[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 在每条连接的 QUIC 握手完成后立即为它构建一个 ConnectionBytes,并与客户端地址一起传给 make_connector 回调,run_hysteria_inbound 正是在这个回调里用 Wire::quic(bytes) 打开会话。bytes() 读取 stats().udp_rx.bytes 和 stats().udp_tx.bytes,这需要获取连接的内部锁。ConnectionBytes 持有连接的一个克隆,因此只要它被持有,就一直能计数;Session、它的注册表条目和它在账本中的记录都通过 Wire 持有它(见谁持有 Wire),并在会话结束时释放它。

supervisor/src/entity/usage.rs
#[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),
}

两个互斥锁都是 parking_lot::Mutex,从不跨 .await 持有。Ledger、Account、Wire 及其方法在 etemenanki_supervisor::entity::usage 中是 pub 的,但没有在 crate 根部重新导出。没有前端程序使用它们;基准测试直接调用 Ledger::bind、take 和 reconcile。entity/user.rs → Account(一对用户名和密码)是另一个同名类型。

存活会话登记所用的 SessionId 来自 Sessions::next,这是一个从 1 开始计数的 AtomicU64,打印为 session#<n>;它在 supervisor 的生命周期内唯一。

方法 锁 作用
Ledger::bind(user, session, wire) -> Arc<Account> 先账本,再该账户 在 accounts 中查找或创建该用户的账户,把 session 以 settled = (0, 0) 插入其 live 表,并把账户插入 active。返回该账户,Session 把它保存在一个 OnceLock 中,以便之后 close。在账本的锁下加入会话,意味着并发的取出要么发现账户已在 active 中,要么发现账户中已带着这个会话。
Ledger::reconcile() 先账本(用于列出 active),再依次每个账户 在账本的锁下把 active 账户克隆到一个 Vec 中,然后释放锁。对每个账户,在其锁下把每个存活会话结算到 pending。
Ledger::take() -> Vec<(UserKey, u64, u64)> 账本,以及依次每个 active 账户 对每个 active 账户:把 pending 移入 reported 并清零;若移出的量不为 (0, 0),就推入 (user, up, down);只有账户仍有存活会话时,才把它留在 active 中。
Ledger::totals() -> Vec<(UserKey, u64, u64)> 账本,以及依次每个账户 对每个曾经绑定过的账户:reported + pending + (wire.read() - settled),最后一项对其所有存活会话求和。不标记任何东西。
Account::close(session) 该账户 把会话从 live 中移除,并把它越过水位的字节加到 pending。

LiveSession::settle 是字节离开 wire、进入账本的唯一位置:

supervisor/src/entity/usage.rs
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
}

take 和 totals 都不保证顺序:两者都遍历一个 HashMap。

supervisor/src/supervisor.rs
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;
}
  • 账本使用 UserKey,前端程序使用 U。UsageBook 通过 names 做转换:names 是一个 parking_lot::RwLock<Vec<U>>,其中 key n 的名字是 names[n - 1](build/users.rs → id_of)。UserKeys 在第一次见到某个 id 时为它分配 key,从 1 开始,并为每个 id 保留它的 key,所以被移除用户的最后一批字节仍以其 id 上报,重新加入的用户会得到同一个 key 和同一个账户。
  • 只有当某个入站准入该用户时才分配 key:build/users.rs → admit 只对入站用户集中持有该协议所读取凭据类型的用户调用 UserKeys::key。仅仅在用户集中还不够,所以没有任何入站准入的用户永远不会得到 key,永远不会绑定会话,也永远不会有账户。
  • Actor::commit_keys 在每次提交(commit)开始时,在存储任何用户表或启动任何监听器之前,把新分配的 key 对应的 id 追加到 names;无论是一次应用(apply)还是一次用户编辑都是如此。key 在准备(prepare)阶段分配在一个暂存副本上,只有变更提交时才保留。因此每个 key 都在会话能够绑定它之前发布。
  • UsageBook::name 查找某个 key,找不到时以 <key> was bound before it was published panic。UserKey 没有实现 Display,所以消息用 {key:?} 格式化 key,例如 UserKey(7) was bound before it was published。已经取出的字节无法放回,所以没有名字的 key 被当作 bug 处理,而不是当作一个可以丢弃的增量。
  • UsageBook::take 调用 Ledger::take,并在读锁下把每个元组映射为 UsageDelta<U>;UsageBook::totals 对 Ledger::totals 和 UsageTotal<U> 做同样的事。
  • take_usage 和 usage_snapshot 是句柄上的普通同步方法。它们直接读取 book,从不经过 actor 的命令通道,所以在应用进行期间以及 shutdown 返回之后都能应答。它们在调用方的线程上运行,并获取账本和各账户的 parking_lot 锁,所以异步调用方在整个遍历期间都占着自己的 worker。
  • Actor 拥有一个 Arc<UsageBook<U>>;start 把它克隆到 Supervisor 句柄和推送任务中。
supervisor/src/serve.rs (drive)
// 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 以 showing_progress 模式运行运行时,所以它产出的每个 Traffic 都是一个增量,所有增量之和等于运行时的总量(服务端运行时)。握手步骤和中继步骤都会运行 count。以错误结束的步骤不产出 Traffic:它在失败之前读写的字节留在运行时尚未交出的增量中(poll_once 只在成功时交出增量),所以 wire 永远看不到它们。drive 用 io::Error::other(e.to_string()) 转换这个错误并返回,serve_connection 以 debug 级别把它记录为 <protocol> connection from <source> ended: <error>,其中源地址是 Debug 形式(Some(<ip>) 或 None)。

supervisor/src/serve.rs
struct WireCounted<S> {
inner: S,
wire: Option<Arc<Wire>>,
}

SOCKS 驱动器(SocksInbound::serve,见 SOCKS)不运行 ProxyServerRuntime,所以没有东西报告它的 stream 搬运了多少字节。serve_connection 改为把 stream 包在 WireCounted 中:poll_read 把本次读取填入的字节数计为上行,poll_write 把它返回的字节数计为下行。flush 和 shutdown 直接透传。同一个调用还把 connector.charging_datagrams() 交给驱动器,它会设置 AppConnector::charge_datagrams。当关联的 FanOutLink 打开一个子 link 时,poll_opening 把 FlowScope::charge(会话的 Arc<Wire>)传给 Tracker::meter_datagram,子 link 的 FlowHandle 就会把它收发的每个载荷字节既计入这个 wire,也计入流自己的计数器。fan-out 见出站。

会话的 Wire 只创建一次,由打开会话的一方创建,并移入 Sessions::open,后者把它包进一个 Arc。之后有三个长期持有者:

持有者 字段 释放时机
Session 句柄 Session::wire,由 Session::wire() 返回给 drive、SOCKS 的 serve_connection,以及计费 UDP 关联的 AppConnector::connect 随最后一个 Arc<Session> 释放
注册表条目 Sessions 中的 Entry::wire,由 Sessions::stats 读取 Sessions::remove,由 Session 的 Drop 调用
用户的账户 LiveSession::wire,来自 Ledger::bind Account::close,由 Session 的 Drop 调用

各连接上的计数方在计数期间都持有它的克隆:drive 的 count 闭包、SOCKS 的 WireCounted::wire,以及 SOCKS UDP 关联中 FanOutLink 里的 FlowScope::charge 和每个计费子 link 中的 FlowHandle::charge。FlowScope::session 只是一个 SessionId,所以 fan-out link 不会让会话保持存活;让会话保持存活的是 SOCKS 驱动器,它通过 AppConnector 做到这一点,并在 SocksInbound::serve 运行期间一直持有它。对 Hysteria 2 会话,持有 Wire 也就持有了它的 quinn::Connection。

第一条流的 Session::admit 在释放注册表的锁之后调用 Ledger::bind,并把返回的账户存入会话的 OnceLock<Arc<Account>>。Session 的 Drop 依次执行三步:取消会话的 token;如果已设置账户,调用 Account::close;把会话从注册表中注销。admit 的各个步骤、让后续流跳过注册表锁的 bound 标志,以及为什么在 Drop 查找账户之前它总是已经设置好,见用户、principal 与会话。

flowchart LR
  RT["drive:Traffic.transport_rx / tx"] --> WC1["Wire::counted"]
  SC["WireCounted:SOCKS stream"] --> WC1
  MD["MeteredDatagram:SOCKS UDP 载荷"] --> WC1
  QS["quinn 统计:udp_rx / udp_tx"] --> CB["ConnectionBytes"]
  CB --> WQ["Wire::quic"]
  WC1 --> B["第一条流时 Ledger::bind"]
  WQ --> B
  B --> ACC["该 UserKey 的 Account"]
  ACC --> TK["take_usage 或 UsageSink"]
sequenceDiagram
  participant D as 数据平面
  participant W as Wire
  participant R as 调和(阻塞线程池)
  participant A as Account
  participant T as 取出
  D->>W: add(up, down)
  R->>A: 加锁,逐个结算存活会话
  A->>W: read()
  Note over A: pending += now - settled, settled = now
  D->>A: Session 被 drop:close(id) 结算并移除它
  T->>A: 加锁:pending 移入 reported
  T-->>T: 每个有应报字节的用户一个 UsageDelta

一个字节恰好经由以下两条路径之一进入 pending,两条路径都在账户的锁下进行:

  1. 调和。 Ledger::reconcile 在每个采样器 tick 运行一次:Sampler::tick 在 spawn_blocking 中首先调用它,tick 间隔为 DEFAULT_SAMPLE_INTERVAL(1 s),除非 SupervisorBuilder::sample_interval 设置了别的值。推送任务在每个批次之前也会调用它。它把每个存活会话自其水位以来搬运的字节移入待取总量,并推进水位。
  2. 关闭。 已绑定会话的最后一个句柄被 drop 时,Drop 调用 Account::close(顺序见谁持有 Wire)。close 在同一个临界区内最后一次结算该会话,并把它从 live 中移除。

随后,一次取出清空 pending。所以一次取出上报的是在它之前已被调和或已关闭的字节:存活会话自上一个 tick 以来搬运的字节由之后的某次取出上报,已结束会话的字节则全部在下一次取出中上报。

stateDiagram-v2
  [*] --> Active: bind(该用户的第一个会话)
  Active --> Active: bind、reconcile、close
  Active --> Kept: take 时没有存活会话
  Kept --> Active: bind(新会话)

处于 Kept 的账户带着其 reported 总量留在 accounts 中,所以 usage_snapshot 仍然列出该用户,但取出不再访问它。它没有任何待取字节:把它移出的那次取出已经清空了它,此后只有绑定才能再向它添加字节(an_idle_reported_account_is_not_visited_until_it_binds_again)。

已绑定会话的 wire 截至会话关闭时统计到的每个字节,都恰好落入一个 take_usage 结果或一个 sink 批次。保证这一点的组成部分:

  • 每个存活会话一个水位。 settled 记录该会话的 wire 中已有多少进入 pending。只有 settle 搬运字节,它在同一步中计算差值并推进水位。
  • 每个账户一把锁。 调和、关闭和取出都只在账户的 Mutex 下操作账户。同一个会话的一次调和与一次关闭不可能结算同一批字节:后运行的一方会看到已推进的水位,或者在 live 中已找不到该会话。
  • 关闭连同会话一起移除水位。 close 把会话并入之后,账本不再引用它的 wire。因此计数覆盖的是会话关闭时 wire 显示的数值。
  • 取出会清零它上报的量。 std::mem::take 清空 pending,并在同一把锁下把相同的量加到 reported。
  • 一致的加锁顺序。 每条同时持有两把锁的路径都先取账本的锁,再取某一个账户的锁(bind、take、totals);reconcile 在获取任何账户的锁之前就释放账本的锁,close 只获取账户的锁。没有任何路径先取账户的锁再取账本的锁。
  • 绑定与取出串行执行。 bind 在账本的锁下把会话插入账户、把账户插入 active,而取出在整个遍历期间都持有这把锁。所以一次取出要么完全发生在绑定之前(会话还不在账本中,其水位将从零开始),要么完全发生在绑定之后(账户已带着该会话位于 active 中)。
  • 只有在不可能欠账时,账户才离开 active。 take 只在清空账户之后、并且 live 为空时,才把账户从 active 中移除。该用户之后的会话会再次绑定,把它放回去。
  • 被移除的用户保留其账户。 移除用户会吊销其 principal,并对会话执行入站的移除策略;账户和 key 保留。被 Keep 策略保留下来的会话继续计在同一用户名下,被移除用户的会话搬运的所有字节,会在被调和或关闭之后由下一次取出上报。
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.
}
  • 每个用户最多一个增量,没有应报字节的用户没有增量,顺序不定。
  • 结果相对存活会话最多滞后一个采样器 tick:存活会话自上一个 tick 以来搬运的字节要等之后的某次取出。已结束会话的字节全部在下一次取出中。
  • 两次互相竞争的取出,或者一次取出与一个 sink 竞争,会在彼此之间分割待取总量:每个字节归先清空其账户的那一方。
  • usage_snapshot() 为每个曾经绑定过会话的用户返回一个 UsageTotal(包括总量为零的用户),根据调用时的水位和存活的 wire 计算。它不标记任何东西,所以无论它是否运行,取出上报的字节都一样。不过和取出一样,totals 在遍历每个曾绑定过的账户期间一直持有账本的锁,并在锁下读取每个存活会话的 wire(读取 Hysteria 2 wire 会获取那条 QUIC 连接的锁)。并发的取出和第一条流的绑定都要等待这次遍历,所以按实际需要的频率调用它即可,不要更频繁。

每个会话的实时 wire 字节(不滞后、不标记)也可以查看:

调用 返回 字段 路径
Supervisor::sessions()(async) Vec<SessionInfo<U>> id、inbound、source、user: Option<U>、up、down、started actor 上的一个命令,它用自己的 UserKeys::id 把每个 UserKey 映射为 U,而不是用 UsageBook::names
Tracker::sessions() Vec<SessionStats> id、inbound、source、user: Option<UserKey>、user_label、up、down、started 直接调用,不经过 actor

两者都来自 Sessions::stats:它在注册表的锁下列出会话,释放锁之后才逐个调用 Wire::read,因为读取 QUIC wire 会获取连接的锁,而打开会话和准入不能等这把锁。up 和 down 是 wire 自会话打开以来的计数,不论会话是否已绑定。注册表见用户、principal 与会话。

SupervisorBuilder::usage_sink(sink, interval) 保存一个装箱的闭包(PushUsage<U>)。SupervisorBuilder::start 应用第一个 spec(期望状态),启动采样器,然后如果设置了 sink,就带着 UsageBook 和 actor 的 background_stop token spawn 推送任务。推送任务的句柄与采样器的句柄一起加入 actor 的 background 列表。

supervisor/src/supervisor.rs (push_usage)
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;
}
}
步骤 细节
第一个批次 start 之后一个 interval:interval_at 不会立即触发。
每个 tick 在阻塞线程池上调和,然后取出,批次非空时调用 report。正是先调和,才使一个批次恰好覆盖它的间隔,而不是止于上一个采样器 tick。
慢 sink 任务先等待 report 完成,再等待下一个 tick。由于使用 MissedTickBehavior::Delay,期间错过的 tick 不会一次性补发;下一个批次包含更多字节。数据平面从不等待 sink(a_slow_usage_sink_does_not_stall_the_data_plane)。
停止 background_stop 被取消时,任务再调和并取出一次,批次非空时上报,然后返回。

想把批次交给另一个任务的 sink 可以通过 channel 转发。report 返回的 future 可以借用批次,但批次只在该 future 完成之前有效,所以要转交批次的 sink 需要复制它:

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.
}
});

Actor::shutdown(grace) 只在每条连接都结束之后才取消 background_stop,所以这些连接的最后一批字节会在推送任务最后一次调和之前进入账本:

sequenceDiagram
  participant C as 调用方
  participant Act as Actor
  participant TT as TaskTracker
  participant Smp as 采样器任务
  participant P as 推送任务
  participant Sk as UsageSink
  C->>Act: shutdown(grace)
  Act->>TT: 停止监听器,取消探测,close,最多等待 grace
  Act->>TT: 取消 root,等待每个任务结束
  Note over TT: 每个结束的 Session 把剩余字节并入账户
  Act->>Smp: 取消 background_stop
  Act->>P: (同一个 token)
  Smp->>Smp: 最后一个 tick:调和,发布
  P->>P: 调和,取出
  P->>Sk: report(最后一个批次),为空则跳过
  Act->>Act: 等待两个任务结束
  Act-->>C: shutdown 返回
  • shutdown 会等待推送任务,所以正在进行的 report future(包括最后一个批次的那一个)会在 shutdown 返回之前完成。与远程服务通信的 sink 应当自己用超时限制这次调用。
  • 如果所有 Supervisor 句柄都在没有调用 shutdown 的情况下被 drop,actor 的命令循环会结束并运行 shutdown(Duration::ZERO),以相同的顺序停止后台任务。
  • take_usage 在 shutdown 之后仍然可用。在 a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown 中,shutdown 之后的 take_usage 为空,因为推送任务的最后一次取出已经清空了账本。
操作 运行时机 访问的对象 持有的锁
Wire::add 每个中继步骤(drive)、每次读写(WireCounted)、每个数据报(计费子 link) 无 无:最多两次 relaxed 原子加法
统计 QUIC 会话 在 quinn 内部 supervisor 一侧无 无
Ledger::bind 每个已绑定会话一次,在 AppConnector::connect 内,运行于准入其第一条流的异步 worker 上 一个账户 先账本,再该账户;取出或 totals 持有账本锁时需要等待
Account::close 每个已绑定会话一次,在会话结束时 一个账户 该账户
Ledger::reconcile 每个采样器 tick 以及每个 sink 批次之前,在阻塞线程池上 每个 active 用户的每个存活会话 账本锁只在列出时持有;每次一个账户
Ledger::take 每次 take_usage 和每个 sink 批次 每个 active 用户(有存活会话或待取字节)一个账户 遍历期间持有账本锁,每次一个账户
Ledger::totals 每次 usage_snapshot 每个曾绑定过的账户,以及每个存活会话的 wire 遍历期间持有账本锁,每次一个账户

取出的开销取决于活跃用户数,而不是会话、流或字节数;调和的开销取决于存活会话数,并且在后台运行。由于取出和 totals 在整个遍历期间都持有账本的锁,而 bind 需要这把锁,所以它们的遍历时长也就是第一条流的准入在异步 worker 上可能等待的时长。

只有 take 会把账户从 active 中移除。一次调和会结算 active 中每个有存活会话的账户;没有存活会话的账户只花费一次加锁和一次空循环。

基准测试 supervisor/benches/tracking.rs 在一个服务器规模的账本上测量一次取出和一次调和(cargo bench -p etemenanki-supervisor --bench tracking)。它的 ledger(per_user) 辅助函数直接在一个 Ledger 上,为 USERS = 50 000 个用户中的每一个绑定 per_user 个会话,编号为 SessionId::new(user * per_user + s),并返回账本、各 wire 和各账户。没有任何会话被关闭,所以每个账户都保有其存活会话,并在每次被测量的取出中始终留在 active 中。

  • take_usage, 50k active users:每个用户分别有 1、10 和 20 个存活会话的账本(50 000、500 000 和 1 000 000 个会话)。每次被测量的取出之前,每个会话的 wire 在每个方向搬运 1 500 字节,账本随后被调和,这是最坏情况。预期各行保持平稳,因为一次取出对每个用户只访问一个账户。
  • reconcile (sampler tick), 50k active users:同样的账本,每次被测量的调和之前,每个会话每个方向搬运 1 500 字节。它的开销随会话数增长。

两组都使用 sample_size(20) 和 BatchSize::PerIteration。该文件的第三组测量被计量流的每字节开销,属于追踪的内容。

内存:一个 Account 包含一个 Arc、一个互斥锁、一个存活会话表和两对 u64;一个 LiveSession 包含一个 Arc<Wire> 和一对 u64,在其会话结束时由 Account::close 移除。

采样器每个 tick 发布一次载荷统计,按用户的用量则单独维护。两者测量的是不同的东西,谁也无法从另一方推导出来:

用量(take_usage、UsageSink、usage_snapshot) StatsSnapshot(Tracker::stats)
字节 会话客户端一侧的线路字节,包括协议开销 每条流与其出站之间往返的载荷(Metered、MeteredDatagram)
计数单位 会话,在第一条流时绑定到一个用户 流:一条 stream,或者 UDP 关联的一个子 link
键 前端程序的 U UserKey,附带用户的标签(UserStats)
匿名流量 任何地方都不统计 计入 total 以及按入站和出站的统计,但不按用户统计
形式 由取出清空的增量,或者累计总量 自启动以来的总量加上上一个 tick 内的速率,每个 tick 在 watch channel 上整体替换
包含哪些用户 每个有应报字节的用户(取出)或每个曾绑定过的用户(snapshot) 有存活流或在上一个 tick 有流量的用户
恰好一次 是,在各次取出与 sink 之间 不是增量:每个 tick 重新发布总量
是否由 REST API 提供 否 是,但不含按用户的行

测试展示了两者的差距。在 mux_sub_flows_and_their_session_count_exactly_what_moved 中,三个 VLESS mux 子流各自恰好统计自己的载荷,而会话和该用户唯一的增量统计的是套接字承载的每个字节(mux.sent 和 mux.received),包括 mux 和 VLESS 头部。在 a_hysteria2_session_bills_its_quic_connections_bytes 中,流在每个方向统计 20 000 字节,用量则更多。SOCKS UDP 关联是两者在数据报上一致的情形:它的子 link 的载荷正是计给会话的字节。

限速跟随载荷一侧:用户的 pacer 按其每条流的载荷扣减(FlowHandle::add_up 和 add_down),而不是按线路字节。用户按线路字节计费,按载荷限速。

不变量 保证机制 固定它的测试
已绑定会话搬运的每个字节都在其用户名下被恰好取出一次(在一个 take_usage 结果或一个 sink 批次中),无论打开、传输、关闭、移除、调和和取出如何交错 账户锁下的水位和 pending every_byte_is_taken_exactly_once
未绑定会话搬运的字节不会被取出 只有带用户 key 的第一条流才调用 bind every_byte_is_taken_exactly_once(其中 10% 的会话保持未绑定)
所有会话都结束并被取出后,不再有剩余,且 totals 等于实际搬运的量 close 并入剩余字节;take 清空 every_byte_is_taken_exactly_once
取出从不上报空增量 take 跳过为 (0, 0) 的应报量 tests/unit/usage.rs 的 sum 辅助函数在每次取出时断言这一点
在其他线程上与会话竞争的取出和调和,仍然把每个字节恰好统计一次 加锁顺序;关闭和调和都在账户的锁下结算 takes_racing_live_sessions_report_every_byte_once
取出只在存活会话的字节被调和之后才上报它们;在此之前 totals 已能显示;关闭会并入剩余部分 take 只清空 pending a_take_reports_live_bytes_once_they_are_reconciled
没有存活会话、也没有待取字节的账户离开 active,并在下一次绑定时回来 take 中的 retain an_idle_reported_account_is_not_visited_until_it_binds_again
移除后得到同一个 key 的用户,沿用同一个账户 Ledger::bind 在 LedgerInner::accounts 中找到它,而后者保留每个账户 every_byte_is_taken_exactly_once(它的 Remove 步骤以新的 principal、同一个 key 重新加入用户)
被移除后重新加入的用户得到同一个 key build/users.rs → UserKeys::key 保留每个 id 的 key 该 proptest 不覆盖:它用固定的 key(UserKey::new(u + 1))构建 principal,从不使用 UserKeys
会话绑定到它准入的第一个 principal Session::admit 的 bound 标志和注册表锁 the_first_flow_binds_the_session_to_its_principal
第一条流出示已吊销 principal 的会话被拒绝,并保持未绑定 Session::admit 在注册表的锁下检查 revoked() a_session_whose_first_flow_presents_a_revoked_principal_is_refused_under_every_policy 固定了 PermissionDenied 拒绝、被取消的 token,以及 stats 中的 user == None。该测试不检查账本:没有字节进入账本,是由 admit 在 Ledger::bind 之前返回推出的
token 已被取消的监听器不打开会话 Sessions::open 在注册表的锁下检查 stop a_stopped_listener_opens_no_session
会话存活期间其 wire 可读,并随最后一个句柄一起消失 Sessions::stats、Session 的 Drop a_session_is_listed_with_its_bytes_until_its_last_handle_drops
流式会话的用量恰好是其套接字在传输层之后承载的字节 drive 计入 transport_rx 和 transport_tx mux_sub_flows_and_their_session_count_exactly_what_moved
SOCKS UDP 载荷在控制流之外另行计给关联的会话 charging_datagrams、FlowScope::charge a_udp_association_counts_each_sub_link_and_charges_its_session
Hysteria 2 会话按其 QUIC 连接的 UDP 字节计费;usage_snapshot 与 take_usage 一致 Wire::quic a_hysteria2_session_bills_its_quic_connections_bytes
通过密码准入的 Hysteria 2 用户在该用户名下计费 来自用户集的 principal a_hysteria2_user_set_admits_by_password
间隔为 200 ms、流量稳定时,批次之间相隔 0.9 到 2 个间隔,连同 shutdown 之后的那一批,总和恰好等于用户搬运的量 push_usage、关停顺序 a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown
慢 sink 从不拖慢连接,也不丢失任何字节 推送任务独立运行;Delay tick a_slow_usage_sink_does_not_stall_the_data_plane
  • 调和中的 panic。 采样器和推送任务都在 spawn_blocking 中运行 reconcile,并用 std::panic::resume_unwind 重新抛出其中的 panic,这会结束该任务。如果采样器的任务以这种方式结束,就不再有任何东西在 tick 上调和存活会话:在没有 sink 的前端程序中,取出只能在存活会话关闭之后才看到它的字节。
  • report 中的 panic。 它会结束推送任务;它持有的批次已被取出,不会再次提供。
  • 未发布的 key。 UsageBook::name 以 <key> was bound before it was published panic。commit_keys 在任何可能让会话绑定新 key 的操作之前运行,所以这表示 supervisor 中存在 bug,而不是一种运行时状况。
  • 会话的取消。 关闭会话(策略、Supervisor::close、被移除的入站、关停)会取消其 token;它的任务结束并 drop 各自的句柄,最后一个句柄的 Drop 把剩余字节并入账户。没有任何路径会在不调用 Account::close 的情况下结束已绑定的会话,因为并入发生在 Drop 中。在第一条流之前被取消的会话永远不会绑定:admit 以 the session was closed 拒绝每一条流。
  • 第一条流拨号失败。 会话在拨号之前就已在 AppConnector::connect 中绑定,所以它保持绑定,其字节计入其用户;失败以拨号失败的形式到达协议核心。
  • 失败的步骤。 以错误结束的那个运行时步骤的字节不会计入 wire(见数据平面在何处给 Wire 计数);之前各步骤统计的字节都保留。
  • accept 中途停止的监听器。 监听器的 token 被取消后,Sessions::open 拒绝打开会话(a_stopped_listener_opens_no_session);恰好在此刻完成握手的 Hysteria 2 连接会得到一个已被取消的 token 和一个没有会话的 connector。两者都不会进入账本。
  • 关停。 推送任务最后的调和与取出在每条连接都结束之后运行;见关停时的推送任务。账本本身不持久化:supervisor 消失之后,新的 supervisor 从零开始。
项 值 定义位置 对用量的影响
DEFAULT_SAMPLE_INTERVAL 1 s,除非 SupervisorBuilder::sample_interval 设置了别的值 track/sampler.rs 存活会话的字节在被取出看到之前最多等待多久
第一个 sink 批次 start 之后一个 interval push_usage(interval_at) 第一个间隔过去之前不推送任何东西
错过的推送 tick MissedTickBehavior::Delay push_usage 迟到的批次覆盖更长的间隔;tick 不会一次性补发
批次大小 每个 active 用户最多一个增量 Ledger::take 随用户数增长,而不是随会话数增长
账户 每个曾绑定过会话的用户一个 LedgerInner::accounts 在 supervisor 的整个生命周期内保留
已发布的名字 每个曾分配的 UserKey 一个 UsageBook::names 在 supervisor 的整个生命周期内保留
计数器 每个会话、每个账户在每个方向一个 u64 Wire、AccountState 统计自会话打开以来(wire)或自 supervisor 启动以来(账户的 reported)的量,每个方向最多 u64::MAX 字节(约 18.4 EB)

supervisor/tests/unit/usage.rs 中的单元测试(作为 tests 模块编译进 entity/usage.rs):

测试 驱动的内容 固定的行为
every_byte_is_taken_exactly_once 一个 proptest,在 4 个用户上执行 1 到 199 个操作:打开会话(权重 2;90% 的情况下绑定到用户)、在存活会话上每个方向传输 0 到 9 999 字节(权重 6)、关闭一个会话(2)、在 UserRemovalPolicy::Close 下移除一个用户并以新的 principal、同一个 key 重新加入(1)、调和(2)、取出(2) 关闭所有会话并做最后一次取出之后,所有增量之和按用户等于已绑定会话搬运的量;没有剩余可取;totals 与之一致
takes_racing_live_sessions_report_every_byte_once 8 个 worker 线程,每个打开 200 个会话,并对每个会话 50 次加上 (3, 7),分布在 3 个用户上;一个线程循环调和、一个线程循环取出,直到各 worker 结束 真实并发下按用户的精确总和
a_take_reports_live_bytes_once_they_are_reconciled 一个会话:add(10, 20)、取出、调和、add(1, 1)、取出、drop、取出 第一次取出为空,而 totals 已显示 (10, 20);第二次上报 (10, 20);关闭把 (1, 1) 并入第三次
an_idle_reported_account_is_not_visited_until_it_binds_again 一个结束并被取出的会话,然后同一用户的一个新会话 账户离开 active;新的绑定使它回来;totals 是两个会话之和

supervisor/tests/tracking.rs 中的端到端测试,每个都在回环套接字上运行真实的 supervisor:

测试 设置 固定的行为
mux_sub_flows_and_their_session_count_exactly_what_moved 一条 VLESS mux 连接,三个子流连到一个 TCP 回显服务;采样器每 50 ms 运行一次 每条流统计自己的载荷;会话统计套接字承载的每个字节;该用户唯一的增量等于这个值,第二次取出为空
a_udp_association_counts_each_sub_link_and_charges_its_session 以 alice 身份建立的 SOCKS UDP 关联,向路由到两个出站的两个 UDP 回显服务发送 4、6 和 5 字节 每个出站一条流;会话和增量为 (26 + 15, 14 + 15):上行为问候 3 字节、认证 13 字节、请求 10 字节,下行为方法 2 字节、状态 2 字节、应答 10 字节,然后每个方向 15 字节载荷
a_hysteria2_session_bills_its_quic_connections_bytes 一个按账户准入 alice 的 Hysteria 2 服务端,由第二个 supervisor 的 Hysteria 2 出站拨号;一次 20 000 字节的回显 流在每个方向恰好统计 20 000;会话和用量统计得更多;usage_snapshot 与 take_usage 一致
a_hysteria2_user_set_admits_by_password 一个使用 Hysteria2UserAuth::Password 的 Hysteria 2 入站,一个密码正确的客户端和一个密码不同的客户端 被准入的用户在其 id 名下计费;另一个客户端被拒绝
a_usage_sink_gets_a_batch_per_interval_and_the_last_at_shutdown 一个 SOCKS 入站,装有一个间隔为 200 ms 的记录用 sink;每 20 ms 做一次 4 字节回显,持续 5.5 个间隔 至少 4 个批次,彼此相隔 0.9 到 2 个间隔;shutdown(Duration::ZERO) 之后各批次之和为 (26 + echoed, 14 + echoed);usage_snapshot 显示同样的总量;take_usage 为空
a_slow_usage_sink_does_not_stall_the_data_plane 同样的入站,一个每次 report 睡眠 1 s 的 sink,间隔 50 ms,持续 1.5 s 的 11 字节回显 每次回显往返都低于 250 ms;在此期间 sink 最多看到 2 个批次;最终总和精确

两个 sink 测试都使用该文件的 Recorder sink 及其 socks_with_sink(sink, interval) 辅助函数,后者启动一个从账户集中准入 alice 的 SOCKS 入站。Recorder::report 在被调用时(在其 future 被 poll 之前)把 (Instant::now(), batch.to_vec()) 推入一个共享的 Vec,并返回 tokio::time::sleep(delay) 作为该 future:第一个测试中是 Duration::ZERO,第二个中是 1 s。Recorder::sum 把记录下的每个增量相加。

supervisor/tests/unit/session.rs 中与用量相关的会话测试: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,以及 a_stopped_listener_opens_no_session(被取消的监听器 token 使 Sessions::open 返回 None,并让 stats 保持为空)。该文件的其余测试见用户、principal 与会话。

没有测试覆盖 UsageBook::name 的 panic、会 panic 的 sink、安装在未经 shutdown 就被 drop 的 supervisor 上的 sink,以及 UserKeys 给重新加入的用户分配同一个 key。修改这些路径时请补上测试。如何运行测试套件,见测试。