流量计费
源码文件:19 个 · 核对版本 katana v3.0.1
katana/src/traffic.rskatana/src/manager/node.rskatana/src/manager/mod.rskatana/src/manager/proxy.rskatana/src/manager/transport.rskatana/src/runtime.rskatana/src/connector.rskatana/src/meter.rskatana/src/serve.rskatana/src/api/mod.rskatana/src/api/newv2board.rskatana/src/api/sspanel.rskatana/src/config.rskatana/tests/unit/traffic.rskatana/tests/unit/connector.rskatana/tests/unit/meter.rskatana/tests/unit/e2e.rskatana/tests/unit/api/newv2board.rskatana/tests/unit/runtime.rs
katana 中每个节点各有一个 NodeTraffic 注册表。它把每个已授权的凭据映射到一个 UserCounter,其中保存该用户的上传和下载总数,以及执行该用户限速的令牌桶。流把字节写入这些计数器;每个轮询周期,节点管理器对它们做一次快照,把总数发给面板,并在面板接受上报后,精确扣减这次上报所携带的量。
本页深入介绍 src/traffic.rs,以及它周边的上报代码:src/manager/node.rs 和两个面板客户端。本页面向修改字节计数方式、修改用户集合变化时计数器如何转移,或修改上报的构建与核对方式的贡献者。同一套行为在运维视角下的说明见流量上报。
src/traffic.rs 负责:
- 计数器。 每个凭据和 uid 一个
UserCounter,包含两个AtomicU64总数和一个TokenBucket。 - 注册表。 当前每个凭据计费到哪个计数器,以及新流拿到哪个计数器(
NodeTraffic::lookup)。 - 用户集合的切换。 两阶段的
prepare和commit:保留未变化用户的计数器,给新用户和速率变化的用户分配新计数器,并把每个被替换下来的计数器暂存起来,使其字节仍会被上报。 - 上报的账目处理。
snapshot读取每个仍有未上报字节的计数器,commit_reported扣减成功上报所携带的量,restore_residuals把失败上报中没有计数器的行交还回去。 - 速率函数。
determine_rate把节点和用户的限制合并为一个令牌桶速率。
它有意不负责:
- 搬运字节和等待令牌桶。这由
src/meter.rs完成(见计量与速率限制)。 - 决定谁可以打开流,以及取消已离开用户的连接。这由
src/connector.rs中的Admission在lookup和PreparedUsers::cancel_keys之上完成(见准入与用户表)。 - 与面板通信。
src/manager/node.rs中的NodeManager::report_traffic通过PanelClient发起上报(见面板客户端)。
NodeTraffic 中的任何内容都不会持久化。注册表的生命周期与其 NodeManager 相同:NodeManager::new 创建它,并在监听器重建以及节点 PanelClient 的每次替换时都保留它。
TokenBucket
Section titled “TokenBucket”一个字节令牌桶,以每秒 rate 字节的速度补充,最多容纳一秒的 rate。rate 为 0 表示不限速:每个方法在碰到锁之前就返回 None。
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>; fn refill(&self, st: &mut BucketState) -> Instant; fn repaid_at(&self, st: &BucketState, now: Instant) -> Option<Instant>;}Mutex 是 parking_lot::Mutex,Instant 是 tokio::time::Instant,因此测试可以用暂停的时间驱动令牌桶。
| 方法 | 作用 |
|---|---|
new |
初始为满:tokens = rate,last = Instant::now()。 |
refill |
补充自 last 以来的 elapsed * rate 个令牌,上限为 rate(一秒的突发量),并把 last 移到当前时刻。 |
charge |
先补充,再全额扣减 n,即使这会让 tokens 变为负数。返回欠额还清的时刻;余额非负时返回 None。 |
ready_at |
先补充,再回答同一个问题,但不扣费。 |
repaid_at |
tokens < 0.0 时返回 now + (-tokens / rate) 秒,否则返回 None。 |
余额允许为负。正是这种欠额让限速成为字节数的属性,而不是 chunk 大小的属性:在 10 000 B/s 的令牌桶上写入 30 000 字节,会消耗 10 000 字节的突发量外加两秒的欠额,该用户任意一条流的下一次传输都要等这两秒过去。如果令牌桶在每一块数据之前等待“足够的令牌”,那么一块超大的数据最多只需等一个补充周期就会被放行。
令牌桶从不阻塞。charge 和 ready_at 只报告一个截止时刻。在 src/meter.rs 中,Gate::poll_open 在每次传输前等待 ready_at 返回 None,传输之后 Gate::sent 和 Gate::received 调用 charge。一个用户所有流的两个方向共用该用户计数器中的同一个令牌桶。
consume(&self, n: usize)(先 charge,再 sleep 到返回的时刻)只在 #[cfg(test)] 下存在。
UserCounter
Section titled “UserCounter”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和add_down是使用Ordering::Relaxed的fetch_add。Gate::sent对用户发出的字节调用add_up,Gate::received对用户收到的字节调用add_down。up和down是 relaxed load。它们是两次独立的读取,不构成一致的一对;也没有任何逻辑依赖它们是一致的,因为之后每个方向都单独提交。commit_reported对每个方向各做一次fetch_sub,减去的恰好是传入的量。它从不写入零。uid和rate不可变。速率变化无法应用到已有的计数器上,因为令牌桶的速率在构造时就已固定;速率变化会替换计数器(见 prepare 与 commit)。
一条流在存活期间,以 Arc<UserCounter> 的形式在其 Gate 中持有计数器。注册表之后正是依靠这个强引用,来判断一个已退出注册表的计数器是否仍可能收到字节。
身份:AuthKey、UserTag、UserEntry
Section titled “身份:AuthKey、UserTag、UserEntry”pub enum AuthKey { Uuid(Uuid), Name(CompactString),}
pub struct UserTag { pub key: AuthKey, pub uid: i64,}
pub struct UserEntry { pub key: AuthKey, pub uid: i64, pub rate: u64,}AuthKey是注册表的键。内核的UserAuthorizationderive 了Eq但没有 deriveHash,所以 katana 使用自己的可哈希键。src/manager/mod.rs中的user_key根据节点类型选择键:V2ray(VMess 或 VLESS)用户以AuthKey::Uuid为键,Trojan、Shadowsocks 或 Hysteria 2 用户以其 email 标签的AuthKey::Name为键(traffic_email,email 为空时退回到 uid 的字符串形式)。UserTag是入站用户表为每个凭据携带的内容,用来代替计数器。流在打开时把自己的标签解析为计数器。用户表可能比替换它的那次刷新活得更久(长连接会保留它认证时所用的表),因此如果把计数器直接写进表里,已离开用户的新流就会继续计费到一个没人上报的计数器上。UserTag::unattributed()是key = AuthKey::Name("")、uid = -1:Shadowsocks 用户表必须有一个单密码槽位,而 katana 从不通过它提供服务,这就是该槽位的占位值。它不应匹配任何已注册用户,因此解析不到计数器。UserEntry是下一个用户集合中的一行。src/manager/mod.rs中的build_user_entries根据面板返回的UserInfo列表构建它们,其中rate: determine_rate(node.speed_limit, u.speed_limit)。在 V2ray 节点上,uuid无法解析的用户没有键,会被跳过并输出警告skipping user {uid}: uuid is not a valid UUID,警告中只写 uid,从不写凭据。以 email 为键的节点类型总能生成键。
determine_rate
Section titled “determine_rate”pub fn determine_rate(node_bps: u64, user_bps: u64) -> u64;取两个非零限制中较小的一个,其中 0 表示该侧不限速:
node_bps |
user_bps |
结果 |
|---|---|---|
0 |
0 |
0(不限速) |
0 |
u |
u |
n |
0 |
n |
n |
u |
min(n, u) |
两个输入的单位都是每秒字节数。面板客户端用 src/api/mod.rs 中的 mbps_to_bps 转换 Mbps:乘以 MBPS_TO_BPS = 1_000_000.0 / 8.0,对小于或等于零的值返回 0。各个面板如何填写这两个限制,见限速与连接限制。
NodeTraffic
Section titled “NodeTraffic”#[derive(Default)]pub struct NodeTraffic { users: RwLock<HashMap<AuthKey, Arc<UserCounter>>>, residuals: Mutex<HashMap<i64, (u64, u64)>>, draining: Mutex<Vec<Arc<UserCounter>>>,}有未上报字节的计数器,恰好位于以下三处之一:
| 字段 | 保存的内容 | 谁仍能写入 |
|---|---|---|
users |
当前注册表:每个凭据一个计数器。新流从这里获取计数器。 | 该凭据的所有存活流,以及任何新流。 |
draining |
已离开注册表(速率变化、用户离开、凭据改绑)但可能仍有写入者的计数器:变化之前打开、仍在收尾的连接。 | 只有已经持有该 Arc 的流。新流无法拿到它们。 |
residuals |
失败上报交还回来的、按 uid 记录的 (up, down) 总数。没有计数器,只有数字。 |
无。 |
各方法按调用方分组:
impl NodeTraffic { pub fn new() -> Self;
// user-set changes pub fn prepare(&self, entries: Vec<UserEntry>) -> PreparedUsers; pub fn commit(&self, prepared: PreparedUsers);
// flows pub fn lookup(&self, key: &AuthKey, uid: i64) -> Option<Arc<UserCounter>>;
// reporting pub fn snapshot(&self) -> Vec<TrafficSnapshot>; pub fn restore_residuals(&self, rows: Vec<(i64, u64, u64)>); pub fn clear_residuals(&self); pub fn prune_draining(&self);
fn add_residuals(&self, rows: Vec<(i64, u64, u64)>);}lookup 只有在已注册的计数器与标签属于同一个 uid 时才返回它:self.users.read().get(key).filter(|c| c.uid == uid).cloned()。已离开的凭据查不到任何东西,现已属于另一个 uid 的凭据也查不到。Admission::admit 把这个 None 变成一次被拒绝的流。
set_users(一次调用完成 commit(prepare(entries)))和 get(按 UserAuthorization 查找,不管 uid)只在 #[cfg(test)] 下存在。
PreparedUsers
Section titled “PreparedUsers”pub struct PreparedUsers { next: HashMap<AuthKey, Arc<UserCounter>>, carried: Vec<(Arc<UserCounter>, Arc<UserCounter>)>, orphaned: Vec<Arc<UserCounter>>, cancel_keys: HashSet<AuthKey>,}
impl PreparedUsers { pub fn cancel_keys(&self) -> &HashSet<AuthKey>;}| 字段 | 内容 |
|---|---|
next |
完整的下一个注册表。 |
carried |
(fresh, old) 对:速率变化迫使同一个 uid 换用新计数器。旧计数器进入排空;它的连接保留。 |
orphaned |
凭据已离开集合或现已属于另一个 uid 的计数器。它们同样进入排空。 |
cancel_keys |
其存活连接必须被取消的键:恰好是 orphaned 的键。速率变化的键不在其中。 |
TrafficSnapshot
Section titled “TrafficSnapshot”pub struct TrafficSnapshot { pub uid: i64, pub up: u64, pub down: u64, pub counter: Option<Arc<UserCounter>>,}counter 决定上报之后这一行如何核对:
Some(counter):字节仍在该计数器中。上报成功时从计数器中扣减;上报失败时原样保留。None:字节已经从NodeTraffic中取出(被取出的残余量,或排空中计数器的最后一行)。上报成功时直接丢弃;上报失败时必须用restore_residuals交还。
面向面板的类型
Section titled “面向面板的类型”pub struct UserTraffic { pub uid: i64, pub upload: i64, pub download: i64,}
impl PanelClient { pub async fn report_user_traffic(&self, t: &[UserTraffic]) -> Result<()>;}PanelClient 分派到所配置面板对应的客户端。Result 是 anyhow::Result。
计数器所在的位置
Section titled “计数器所在的位置”flowchart LR P["面板下发的用户集合"] --> PR["prepare"] PR --> CM["commit"] CM --> U["users(注册表)"] CM --> D["draining"] F["新流"] -->|lookup| U G["存活流的 Gate"] -->|"add_up / add_down"| C["Arc UserCounter"] U -. 持有 .-> C D -. 持有 .-> C S["snapshot"] --> U S --> D S --> R["residuals"] X["失败的上报"] -->|restore_residuals| R
prepare 与 commit
Section titled “prepare 与 commit”一次用户集合变化分两次调用。prepare 构建下一个注册表,不修改任何状态,也不读取任何字节总数;commit 把它发布出去。在两者之间,调用方构建下一批用户表;构建出错时可以直接丢弃 PreparedUsers,不改变任何东西。
prepare
Section titled “prepare”prepare 在整个执行期间持有 users.read(),并对每个 UserEntry 做出判断:
flowchart TB
E["UserEntry key, uid, rate"] --> Q{"key 在注册表中?"}
Q -->|否| N["新计数器"]
Q -->|是| S{"uid 和 rate 都相同?"}
S -->|是| K["保留现有的 Arc"]
S -->|否| UID{"uid 相同?"}
UID -->|"是:速率变化"| CA["新计数器,旧的进入 carried"]
UID -->|"否:改绑"| OR["新计数器,旧的进入 orphaned,key 进入 cancel_keys"]
之后,注册表中不在 next 里的每个键都是已离开的用户:其计数器进入 orphaned,其键进入 cancel_keys。
| 变化 | 计数器 | 旧计数器 | 存活连接 |
|---|---|---|---|
| key、uid 和 rate 都不变 | 同一个 Arc |
无 | 保留,计量不中断 |
| 新的 key | 新建 | 无 | 尚无 |
| key 和 uid 相同,rate 变化 | 新建,使用新速率 | 进入 carried,排空 |
保留,仍计费到旧计数器并受其限速 |
| key 相同,uid 不同(改绑) | 新建,归属新 uid | 进入 orphaned,排空,仍以旧 uid 上报 |
取消 |
| key 不再出现在列表中 | 无 | 进入 orphaned,排空 |
取消 |
因此,速率变化只对新流生效。变化之前打开的连接继续使用旧计数器及其旧令牌桶,直到连接关闭。两个计数器都以同一个 uid 上报,report_traffic 会把它们合并为一行。
commit
Section titled “commit”pub fn commit(&self, prepared: PreparedUsers) { let PreparedUsers { next, carried, orphaned, cancel_keys: _, } = prepared; let mut users = self.users.write(); let mut draining = self.draining.lock(); draining.extend(carried.into_iter().map(|(_fresh, old)| old)); draining.extend(orphaned); *users = next;}commit 同时持有两把锁,获取顺序与 snapshot 相同,并在替换注册表的同一个临界区内,把每个被替换下来的计数器移入 draining。因此快照看到的每个被替换计数器都恰好只在一处。如果某个计数器可能同时被看到仍在注册表中又已经在排空,它的字节就会被上报两次,随后被扣减两次。
两个阶段都不读取字节总数,所以都不需要等待旧 generation(一代实例)的连接停止。这些连接在 commit 之后写入的内容,都落进一个 snapshot 会持续读取的排空中计数器。
所有调用方都运行在节点自己的任务(NodeManager::run)上,因此用户集合变化和上报是串行的。
| 调用方 | 时机 | 顺序 |
|---|---|---|
ProxyManager::refresh(src/manager/proxy.rs) |
用户集合或节点限速发生变化,且保留监听器。 | prepare,构建替换用的用户表(出错则返回,不改变任何东西),Admission::commit,然后替换用户表。 |
TransportManager::start(src/manager/transport.rs) |
构建监听器时:启动时以及每次重建。 | prepare,构建并绑定所有可能失败的部分,然后 traffic.commit。此时还没有任何服务在运行,所以不需要取消任何连接。 |
NodeManager::reconcile 和 NodeManager::bring_up(src/manager/node.rs) |
面板返回了空的用户列表。 | commit(prepare(Vec::new())):所有计数器进入排空。 |
Admission::commit 持有准入的 leases 锁,调用 NodeTraffic::commit,之后才取消 cancel_keys 中每个键的租约。由于 Admission::admit 在同一把 leases 锁下调用 lookup,一条流要么在 commit 之前被准入(随后被这次 commit 取消),要么按 commit 发布的注册表接受检查。注册表先于新用户表发布,所以新加入的用户不会出现“已通过新表认证,却被尚不认识他的注册表拒绝”的情况。详见准入与用户表。
计数器的生命周期
Section titled “计数器的生命周期”stateDiagram-v2 [*] --> Registered: prepare 创建,commit 发布 Registered --> Registered: 上报成功,commit_reported Registered --> Draining: 速率变化、用户离开或改绑 Draining --> Draining: 仍有写入者持有时 snapshot Draining --> FinalRow: snapshot 看到 strong_count 为 1 FinalRow --> [*]: 上报成功 FinalRow --> Residual: 上报失败,restore_residuals Residual --> Reported: 下一次 snapshot 取出 map Reported --> [*]: 上报成功 Reported --> Residual: 上报失败
已注册或排空中的计数器从不被清零:上报只会扣减它所携带的量。一旦最后一条流释放了它的 Arc,就只剩 draining 向量持有该计数器,没有任何途径能再写入它,它的下一行快照就是最后一行。此后这些字节只以数字的形式存在于快照行或 residuals 中,直到某次上报成功。
snapshot
Section titled “snapshot”snapshot 为每个可能有未上报字节的计数器生成一行,并为每条暂存的残余量生成一行。
- 先获取
users.read(),再获取draining.lock(),同时持有两者。 - 对每个已注册的计数器,推入一行,包含其当前的
up、down以及counter: Some(clone)。 - 对每个排空中的计数器:
- 先于字节读取
Arc::strong_count。如果为1,说明只有向量持有该计数器,已没有写入者,而且这个计数不可能再增加。 - 如果为
1,执行std::sync::atomic::fence(Ordering::Acquire)。该 fence 与最后一个写入者Arcdrop 时的 release 递减配对,使该写入者加上的每个字节对随后的读取都可见。 - 读取
uid、up和down。 - 如果它是唯一持有者,推入
counter: None的一行,并把该计数器从draining中swap_remove。这一行就是最后一行。 - 否则推入
counter: Some(clone)的一行,并保留该计数器。
- 先于字节读取
- 释放两把锁。
- 获取
residuals.lock(),把整个 mapdrain()为counter: None的行。
结果不做合并:同一个 uid 可以出现多次,例如一个存活计数器、一个速率变化后排空中的计数器,以及一条来自之前失败上报的残余量。零字节的行也会包含在内;report_traffic 会跳过它们。
| 锁 | 类型 | 获取者 |
|---|---|---|
NodeTraffic::users |
parking_lot::RwLock |
prepare(读)、commit(写)、lookup(读)、snapshot(读) |
NodeTraffic::draining |
parking_lot::Mutex |
commit、snapshot、prune_draining |
NodeTraffic::residuals |
parking_lot::Mutex |
add_residuals、restore_residuals、clear_residuals、snapshot |
Admission::leases |
parking_lot::Mutex |
Admission::admit、Admission::commit、Admission::retire_all |
顺序是 leases,然后 users,然后 draining。residuals 从不与其中任何一把同时持有。任何 NodeTraffic 或 Admission 的锁都不会跨 .await 持有:上报请求运行时不持有任何一把锁,同时流继续向原子变量写入。
NodeManager::poll_cycle 获取节点信息和用户,执行 reconcile,刷新规则,然后调用 report_traffic,再调用 report_illegal。定时器每 update_periodic 秒运行一次该周期(default_update_periodic 返回 60;poll_period 会应用 .max(1));配置修改和关闭还会触发额外的运行,见何时运行上报。
sequenceDiagram
participant N as NodeManager 任务
participant T as NodeTraffic
participant F as 流(Gate)
participant P as PanelClient
participant S as 面板
N->>T: snapshot()
T-->>N: 行(counter 为 Some 或 None)
N->>N: 跳过零字节行,按 uid 合并,拆分为 commits 和 residuals
F->>T: add_up / add_down(持续进行)
N->>P: report_user_traffic(reports)
P->>S: POST
alt Ok
S-->>P: 已接受
P-->>N: Ok(())
N->>T: 对每个 counter 行 commit_reported(up, down)
else Err(发送错误、超时、HTTP 4xx 或 5xx、SSPanel 响应体或 ret)
P-->>N: Err
N->>T: restore_residuals(没有 counter 的行)
N->>N: warn "report traffic"
end
何时运行上报
Section titled “何时运行上报”| 触发条件 | 运行的内容 | 位置 |
|---|---|---|
定时器触发,每 update_periodic 秒一次。第一次触发发生在节点首次启动完成后一个完整间隔。 |
poll_cycle(false) |
NodeManager::serve |
节点启动完成后的一次配置修改,它换入新的 PanelClient(任何保持身份元组不变的 [node.api] 修改,或只改变字母大小写的 panel_type 修改),或迫使监听器重建(route、listen_ip、cert、disable_sniffing、enable_vless、[node.hysteria]) |
立即运行 poll_cycle(rebuild),其中包含上报 |
NodeManager::apply_static |
| 节点关闭、被移除或被重新创建 | tear_down,然后 report_traffic 和 report_illegal,各尝试一次 |
NodeManager::run |
由配置修改触发的上报就是一次普通上报:节点保留它的 NodeTraffic,而 report_traffic 在运行时通过 self.api() 读取客户端,所以修改之后的上报经由新客户端发出。修改不会移动定时器,除非它同时改了 update_periodic:这时 serve 会开始一个新的间隔,下一次触发发生在一个完整的新周期之后。节点仍在启动、其 bootstrap 仍在重试期间,不会运行任何上报,此时的配置修改也不触发轮询:节点保存它,若它无法构建则拒绝它,然后立即开始下一次尝试。
report_traffic
Section titled “report_traffic”src/manager/node.rs 中的 report_traffic(&self),逐步说明:
- 如果设置了
[node.controller].disable_upload_traffic:clear_residuals()、prune_draining(),然后返回。见禁用上报。 snapshot()。- 对每个至少有一个方向非零的行:
- 以 uid 为键,用
saturating_add把它累加进totals: HashMap<i64, (u64, u64)>; - 如果
counter是Some,把(counter, up, down)推入commits; - 如果
counter是None,把(uid, up, down)推入residuals。
- 以 uid 为键,用
- 如果
totals为空,不发请求直接返回。 - 把
totals转换为Vec<UserTraffic>,把每个u64转换为i64。 self.api.report_user_traffic(&reports).await。- 返回
Ok(())时,对commits中的每一项调用counter.commit_reported(up, down),参数是该行携带的量。 - 返回
Err(e)时,调用self.traffic.restore_residuals(residuals),并以WARN级别记录node {id}: report traffic: {e}。counter 行不需要任何处理:它们的字节从未被取出。
这里有意维护两个独立的视图:
- 线上视图按 uid 合并,因为 newV2board 的负载是一个以 uid 为键的 JSON 对象,第二行会悄无声息地覆盖第一行。
- 提交视图每个计数器一项。一个同时有存活计数器和排空中计数器的 uid 会得到两次
commit_reported调用,每次扣减从对应计数器读出的量。
commit_reported 作用于 Arc 而不是注册表键,所以无论那时计数器位于何处,扣减都会落在正确的计数器上。
两个周期的完整示例
Section titled “两个周期的完整示例”uid 7 有一个已注册的计数器,一条流持续写入,第一次上报失败,第二次成功。图中数字是该计数器的 up 总数。
sequenceDiagram participant F as 流(Gate) participant C as UserCounter uid 7 participant N as report_traffic participant P as 面板 F->>C: add_up(1000),up = 1000 N->>C: snapshot 读到 1000 F->>C: add_up(200),up = 1200 N->>P: push uid 7 upload 1000 P-->>N: Err Note over N,C: counter 行,无需交还,up 仍为 1200 N->>C: 下一次 snapshot 读到 1200 N->>P: push uid 7 upload 1200 F->>C: add_up(50),up = 1250 P-->>N: Ok N->>C: commit_reported(1200, 0),up = 50 Note over C: 这 50 字节等待下一个周期
失败的上报没有造成任何损失:它的 1000 字节从未从计数器中取出,所以第二次快照把它们和期间新到的 200 字节一起带上。第二次请求进行期间写入的 50 字节在提交后得以保留,因为提交扣减的是 1200,而不是写入零。对于没有计数器的行(残余量,或排空中计数器的最后一行),保住字节的是失败分支:restore_residuals 把它们重新暂存,下一次 snapshot 再把它们取出放进结果。
实际发送的内容
Section titled “实际发送的内容”newv2board::Client::report_user_traffic 构建一个从 uid 到 [upload, download] 的 HashMap<String, [i64; 2]>,以 JSON 形式发送到 PUSH_PATH = "/api/v1/server/UniProxy/push",查询字符串中带有 node_id、node_type 和 token。
{ "1001": [52428800, 1073741824] }只要状态码能通过 error_for_status(低于 400)即视为成功。响应体不会被读取。
sspanel::Client::report_user_traffic 把每一行映射为 TrafficItem { user_id, u, d },并通过 post_data 把 PostData { data } POST 到 /mod_mu/users/traffic,查询字符串中带有 key、muKey 和 node_id。
{ "data": [{ "user_id": 1001, "u": 52428800, "d": 1073741824 }] }成功需要满足:状态码低于 400,响应体能被完整读取,并且当响应体能解析为 Envelope(ret 和 data,两者都是 #[serde(default)])时,ret == 1。由于 ret 默认为 0,不含 ret 的 JSON 对象会失败;而完全无法解析为 envelope 的响应体算作成功。ret 不为 1 时以 /mod_mu/users/traffic: panel returned ret={ret} 失败。
两个客户端都使用节点共享的 reqwest::Client,由 build_http_client(api.timeout_secs()) 构建,超时作用于整个请求。ApiConfig::timeout_secs 返回 [node.api].timeout,为 0 时返回 5。修改 timeout 会为同一个节点和同一个 NodeTraffic 换入新的客户端,因此该修改触发的那次上报以及之后的每次上报都使用新的超时。error_for_status 会从错误中去掉 URL,因为面板密钥放在查询字符串里。
| 不变量 | 机制 | 由谁固定 |
|---|---|---|
| 上报失败不丢字节。 | counter 行只会被 commit_reported 减少,而它只在 Ok 时运行。没有计数器的行在 Err 时由 restore_residuals 交还。 |
注册表这一半:restored_residuals_are_retried(tests/unit/traffic.rs)。report_traffic 中在两种结果之间选择的分支没有测试。 |
| 上报进行期间到达的字节得以保留。 | commit_reported 是按上报量做的 fetch_sub,从不写入零。 |
commit_reported_preserves_concurrent、a_departed_users_late_bytes_are_still_reported(tests/unit/traffic.rs) |
| 同一周期内没有计数器被上报或提交两次。 | commit 和 snapshot 以相同顺序同时持有 users 和 draining,因此被替换下来的计数器恰好只在其中之一。 |
rate_change_drains_old_counter_and_reports_once(tests/unit/traffic.rs),顺序执行。没有测试让 commit 与 snapshot 竞争。 |
| 排空中计数器的字节会一直上报到其最后一个写入者消失,且最后一行只上报一次。 | 先于字节读取 strong_count == 1、Acquire fence,以及在同一临界区内 swap_remove。 |
draining_counter_with_live_writer_is_retained、a_departed_users_late_bytes_are_still_reported(tests/unit/traffic.rs) |
| 已离开用户的字节仍会上报,且只上报一次。 | prepare 把计数器移到 orphaned,commit 移到 draining,snapshot 输出最后一行并丢弃该计数器;residuals 由 snapshot 取出。 |
dropped_user_bytes_become_residuals(tests/unit/traffic.rs) |
| 字节归属产生它们的 uid。 | 改绑凭据的旧计数器被置为 orphaned,而不是被复用;lookup 按 uid 过滤,所以被旧连接钉住的用户表无法计费到新账户。 |
rebound_credential_reports_the_old_uid_separately(tests/unit/traffic.rs)、a_credential_rebound_to_another_uid_is_refused(tests/unit/connector.rs) |
| 未变化的用户在刷新和重建之间保持同一个计数器。 | uid 和 rate 都匹配时,prepare 复用现有的 Arc;NodeTraffic 比每个 TransportManager 都活得久,也比每次 PanelClient 替换都活得久。 |
rate_change_drains_old_counter_and_reports_once(前半部分);用户刷新由 unchanged_user_survives_user_refresh(tests/unit/e2e.rs)覆盖。监听器重建时字节的连续性没有测试。 |
| 速率变化会作用于新流。 | 新计数器带有按新速率创建的新 TokenBucket;lookup 把它返回给每条新流。 |
rate_change_drains_old_counter_and_reports_once |
| 已离开的用户无法打开新流。 | lookup 返回 None;Admission::commit 在发布注册表之后取消该用户的租约。 |
a_user_the_registry_does_not_know_is_refused、the_lease_reaches_the_connection_and_goes_with_the_user(tests/unit/connector.rs) |
| 账目状态有上限。 | add_residuals 用 saturating_add 按 uid 合并并跳过全零行,所以无论用户如何频繁变动,residuals 每个 uid 最多一项;排空中的计数器在最后一次快照时离开向量。 |
清理:draining_counter_with_live_writer_is_retained、rate_change_drains_old_counter_and_reports_once。add_residuals 中按 uid 的合并没有专门的测试。 |
| 上报不会在快照与结果之间被放弃。 | 在 NodeManager::serve 中,无论是定时器触发还是配置修改,poll_cycle 都在 tokio::select! 的分支处理体中运行,而处理体运行期间不会轮询关闭分支,所以被取消的节点会先完成当前周期。bootstrap 与关闭竞争的那次启动尝试从不上报。 |
没有专门的测试固定。 |
由此还能推出两个性质:
- 计数器永远不会低于零。 唯一的减法是
commit_reported,减去的是之前从同一个计数器读出的量,而其他所有操作都只做加法。对同一个计数器上报两次会破坏这一点,这也是“只在一处”这条不变量重要的另一个原因。 - 空闲用户不产生任何网络开销。 零字节行会被跳过,没有内容可发的周期不会发出请求。
失败路径与取消
Section titled “失败路径与取消”report_user_traffic 在以下情况返回 Err:连接或发送错误、请求超时、状态码为 400 或以上,以及在 SSPanel 上响应体无法读取或 envelope 的 ret 不为 1。katana 不会在周期内重试,也没有退避:下一次尝试就是下一个 poll_cycle,其快照包含同样的字节以及期间新到的字节。一次长时间中断最终会变成一次更大的上报。
关闭与节点替换
Section titled “关闭与节点替换”节点的 CancellationToken 触发后,NodeManager::run 退出轮询循环,或停止重试 bootstrap,调用 tear_down(关闭 TransportManager:停止接受连接、让每个用户退役、释放监听器),然后再运行一次 report_traffic 和 report_illegal。在监听器绑定之前就被停止的节点没有统计过任何字节,所以这最后一次 report_traffic 找不到任何内容,也不会发出请求。
当一次重载移除某个节点或改变其身份元组(panel_type、api.host、node_id、api.key、面板节点类型)时,走的也是这条路径:apply_reload 会取消旧节点的 token 并等待其任务结束,然后才启动带有全新空 NodeTraffic 的新 NodeManager。面板节点类型来自 src/api/mod.rs 中的 panel_node_type:
- 在 newV2board 上,它是 UniProxy 请求所携带的
node_type:设置了enable_vless且node_type为V2ray、Vmess或Vless(不区分大小写)时为vless,否则为小写的node_type。UniProxy 按 id 和这个类型查找节点,所以改变这个值的修改指向的是另一个面板节点:不只是大小写的node_type修改(在设置了enable_vless时于 V2ray 系名称之间切换除外,因为它们都以vless请求),或在V2ray、Vmess节点上切换enable_vless。katana 会用全新的NodeTraffic重新创建该节点,而不是把一个面板节点的计数器计费到另一个面板节点上。旧节点的最后一次上报经由它的旧客户端发出,发往这些字节当初计入的那个面板节点。 - 在 SSPanel 上,它为空,因为
mod_mu只按 id 查找节点。
其他所有修改,包括 api.timeout,都作为配置更新送达正在运行的节点。节点保留它的 NodeTraffic;对 [node.api] 的修改会换入新的 PanelClient 并立即运行一次上报(见何时运行上报)。只要有任何一个节点构建失败(例如 node_type 未知),整次重载就会被拒绝,因此不会有正在运行的节点因它而被停止或替换。重载这一侧的内容见进程运行时与重载。
最后这次上报只尝试一次。如果失败,它的字节会随旧的 NodeTraffic 一起消失。崩溃或 SIGKILL 会丢失所有尚未上报的内容,因为没有任何东西写入磁盘。
reconcile 拆除监听器并调用 commit(prepare(Vec::new()))。所有已注册的计数器都移入 draining;该周期自己的 report_traffic 随后上报它们,它们的连接全部消失后便离开向量。bring_up 在被要求启动一个没有用户的节点时也会这样做。
设置 disable_upload_traffic 后,report_traffic 从不调用 snapshot:
clear_residuals在每个周期清空残余量 map。prune_draining只保留strong_count > 1的排空中计数器,其余计数器的字节被丢弃。- 已注册的计数器既不读取也不扣减。它们持续累加,其令牌桶也持续限速。
该标志在每个周期都会读取,并可以通过一次身份不变的重载修改。把它改回 false 后,下一次快照会上报已注册计数器在此期间累积的全部字节。
取消用户的租约会在两处结束其连接:租约取消后,每次 Gate::poll_open 都返回 ConnectionAborted(“the user was retired”);在流式监听器上,连接任务会监视其第一条被准入的流所发布的租约(src/serve.rs 中的 until_retired),并结束整个连接。计费不会等待这些:流结束前搬运的字节已经在计数器中,而计数器最后一个 Arc 的 drop 使下一次快照能输出其最后一行。
| 项目 | 值 | 位置 |
|---|---|---|
| 令牌桶容量(突发) | 一秒的 rate:min(tokens + elapsed * rate, rate) |
TokenBucket::refill |
| 初始余额 | 满:rate 个令牌 |
TokenBucket::new |
| 不限速 | rate == 0 |
TokenBucket::charge、TokenBucket::ready_at、determine_rate |
| Mbps 换算 | MBPS_TO_BPS = 1_000_000.0 / 8.0,十进制兆比特 |
src/api/mod.rs |
| 上报间隔 | update_periodic,默认 60 秒,至少 1 秒 |
default_update_periodic、NodeManager::poll_period |
| 首次上报 | 节点首次启动完成后一个完整间隔;第一次 interval.tick() 在进入循环前被消耗。替换面板客户端或迫使监听器重建的配置修改会立即上报。 |
NodeManager::serve、NodeManager::apply_static |
| 上报超时 | [node.api].timeout,为 0 时取 5 秒,作用于整个请求;修改从下一次上报起生效 |
ApiConfig::timeout_secs、build_http_client |
| 重试 | 周期内不重试;由下一个周期带上这些字节 | NodeManager::report_traffic |
| 残余量 map 大小 | 每个 uid 最多一项 | NodeTraffic::add_residuals |
| 计数器宽度 | 每个方向 u64;合并后的总数使用 saturating_add;以 i64 发送 |
UserCounter、report_traffic、UserTraffic |
| 持久化 | 无 | NodeManager::new |
tests/unit/traffic.rs(作为 tests 模块被包含进 src/traffic.rs):
| 测试 | 固定的行为 |
|---|---|
determine_rate_min_nonzero |
determine_rate 的四种情况。 |
rate_change_drains_old_counter_and_reports_once |
uid 和 rate 相同时保留计数器及其字节;速率变化得到一个按新速率创建的空计数器;旧字节只上报一次,排空完的计数器随之消失。 |
draining_counter_with_live_writer_is_retained |
有写入者的排空中计数器以 counter: Some 上报并继续累加,写入者释放后变成一行最终的 counter: None。 |
dropped_user_bytes_become_residuals |
已离开用户的字节以没有计数器的行出现一次,之后不再出现。 |
a_departed_users_late_bytes_are_still_reported |
移除用户的 commit 之后写入的字节仍会上报;一次提交之后,最后一行只携带剩余部分。 |
restored_residuals_are_retried |
交还的行会出现在下一次快照中;clear_residuals 清空 map。 |
rebound_credential_reports_the_old_uid_separately |
旧 uid 保留自己的字节;新 uid 从零开始。 |
set_users_drops_absent |
未列出的键不再能解析。 |
commit_reported_preserves_concurrent |
读取与提交之间的增量得以保留。 |
token_bucket_unlimited_is_instant |
rate == 0 从不等待。 |
token_bucket_rate_limits |
突发量用完后,以 10 000 B/s 传输 10 000 字节至少耗时 0.8 秒真实时间。 |
a_chunk_larger_than_the_burst_is_still_limited |
以 10 000 B/s 扣费 30 000 字节恰好耗时两秒(暂停的时间)。 |
debt_holds_back_the_next_charge_too |
一次 25 000 字节的扣费留下 1.5 秒欠额,下一次 1 字节的扣费要等它还清。 |
an_idle_bucket_banks_one_second_and_no_more |
空闲 60 秒后,一次 20 000 字节的扣费仍会留下一秒欠额。 |
相关的其他测试:
| 文件 | 测试 | 固定的行为 |
|---|---|---|
tests/unit/meter.rs |
the_limit_holds_however_the_writes_are_sized、the_limit_is_shared_by_both_directions、each_direction_is_billed_to_the_user |
通过 Metered 观察到的令牌桶欠额,以及每个字节计费到哪个方向。 |
tests/unit/connector.rs |
a_user_the_registry_does_not_know_is_refused、a_credential_rebound_to_another_uid_is_refused、an_admitted_stream_is_billed_to_its_user、the_lease_reaches_the_connection_and_goes_with_the_user、udp_is_billed_after_routing_and_blocked_packets_are_free |
lookup 作为准入检查,以及计费的字节落入已注册的计数器。 |
tests/unit/e2e.rs |
vmess_traffic_is_metered_and_reported、unchanged_user_survives_user_refresh、proxy_outbound_relays_and_meters、a_hysteria_node_relays_and_meters |
针对一个记录每次 push 请求体的模拟 UniProxy 面板测试完整路径:至少中继的字节会以正确的 uid 到达面板,覆盖用户刷新和用户移除。这些断言都是下限,所以不能证明恰好上报一次。 |
tests/unit/runtime.rs |
an_sspanel_node_is_its_panel_node_id、a_newv2board_node_is_also_the_type_it_asks_for、a_newv2board_type_edit_respawns_the_node |
哪些修改会改变身份元组,从而用全新的 NodeTraffic 重新创建节点:在 newV2board 上,把 node_type 从 V2ray 改为 Trojan,或在 V2ray 节点上切换 enable_vless,会;只改大小写的 node_type 修改、Trojan 节点上的 enable_vless,以及两种面板上的 api.timeout 修改,都不会。 |
tests/unit/api/newv2board.rs |
push_body_shape |
HashMap<String, [i64; 2]> 序列化为 {"uid": [upload, download]} 的形状。它自己构建 map,并不调用 report_user_traffic。 |
report_traffic 自身的分支(按 uid 合并、拆分为 commits 和 residuals、失败分支)只由端到端测试覆盖,而这些测试只走成功路径。修改该函数时应附带一个测试:让模拟面板使一次 push 失败,并检查下一次 push 携带了同样的字节。如何运行各测试套件见测试。