跳转到内容

流跟踪、统计与限速

源码文件:47 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
  • Etemenanki/supervisor/src/track/mod.rs
  • Etemenanki/supervisor/src/track/metered.rs
  • Etemenanki/supervisor/src/track/pace.rs
  • Etemenanki/supervisor/src/track/sampler.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/topology/outbound/udp_fanout.rs
  • Etemenanki/supervisor/src/topology/outbound/dns.rs
  • Etemenanki/supervisor/src/topology/balancer.rs
  • Etemenanki/supervisor/src/topology/plane.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/entity/user.rs
  • Etemenanki/supervisor/src/entity/session.rs
  • Etemenanki/supervisor/src/entity/usage.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/build/inbound.rs
  • Etemenanki/supervisor/src/topology/inbound/mod.rs
  • Etemenanki/supervisor/src/lib.rs
  • Etemenanki/supervisor/Cargo.toml
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/relay.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/mux/mod.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/app/src/lower.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/webclient/src/router.rs
  • Etemenanki/webclient/src/tracking.rs
  • Etemenanki/ffi/src/proxy.rs
  • Etemenanki/supervisor/tests/unit/track.rs
  • Etemenanki/supervisor/tests/tracking.rs
  • Etemenanki/supervisor/tests/hot_swap.rs
  • Etemenanki/supervisor/benches/tracking.rs
  • Etemenanki/webclient/tests/tracking.rs
  • katana/src/lower/mod.rs
  • katana/src/manager/node.rs

supervisor(监管器)会跟踪它路由的每一条流。一条流就是一个经过路由的出站:一条 stream,或者 UDP 关联中通往某个出站版本的一条子链路。一条普通的代理 TCP 连接是一条流;一条 mux 连接或 Hysteria 2 连接会打开多条流,gRPC 连接也是如此,它的每条 HTTP/2 stream 都是一个独立的会话。connector(连接器)把它打开的每个出站包装成 Metered stream 或 MeteredDatagram link。包装器把流登记到一个并发 map 中,统计它两个方向的有效载荷,在它被强制关闭(kill)后让它失败,并在其用户的令牌桶处于欠额时对它节流。一个采样器任务每个 tick(采样周期)读取一次所有流的计数器,并按入站、出站和用户发布累计值与速率。

本页面向修改 supervisor/src/track/、connector 或 UDP fan-out 的贡献者,以及读取 Tracker 的前端程序作者。内容包括注册表、包装器、强制关闭开关、生命周期事件、采样器与 StatsSnapshot、每用户限速,以及注册表设计背后的基准测试。每个会话统计的线上字节(wire bytes),以及采样器调和的用量账本,见每用户用量统计。流如何被路由,以及每个包装器内部的排空守卫,见 plane(数据平面):为每条流路由。

跟踪代码(supervisor/src/track/)负责:

  • 登记 connector 打开的每条流,附带在打开时固定下来的元数据(FlowMeta),并在其包装器被 drop 时注销它;
  • 在只属于该流的计数器上统计每条流的有效载荷;
  • 从中继该流的任务之外关闭一条流,或关闭某个谓词选中的所有流;
  • 在一个有界 broadcast 通道上公布每次打开和关闭;
  • 每个 tick 对所有流采样一次,生成 StatsSnapshot 并发布到 watch 通道上,并在同一个 tick 里调和用量账本;
  • 用被限速用户唯一的令牌桶对其流节流,这个令牌桶由该用户的所有流在两个方向上共享。

它不负责:

  • 统计线上字节。这由会话的 Wire 负责,再由账本计费(每用户用量统计)。流计数器统计的是有效载荷。唯一的例外是计入会话的数据报流(SOCKS UDP ASSOCIATE 的子链路),它的包装器还会把有效载荷加到所属会话的 Wire 上;见流在哪里被计量中的 charge 一列。
  • 关闭会话。Supervisor::close 通过 actor 按会话、用户、入站或全部关闭整个会话(用户、principal(身份主体)与会话)。强制关闭只结束一条流,会话及其其他流继续运行。
  • 路由或排空。流在被包装之前就已完成路由;一旦排空策略关闭了流所在的出站版本,包装器内部的 Guarded link 就以 the outbound this flow was opened on was drained 让流失败(plane)。
  • 看到 supervisor 为自己拨号的连接。RoutedDialer(supervisor/src/topology/outbound/dns.rs)为 DNS 解析器到其服务器的连接做路由,并直接用 Target::connect_stream 打开它们;负载均衡器的健康探测(supervisor/src/topology/balancer.rs → probe)是一次裸 TCP 连接。两者都不经过 connector,所以都不会在这里登记或计数。

模块文档把它的设计规则表述为数据平面从不等待观察者,并用两条性质来定义:

  • 每条流的计数器归它自己所有,只由中继该流的任务写入,因此没有争用。没有任何按 tag 或按用户的计数是逐字节进行的;采样器每个 tick 汇总一次各流的计数器。
  • 生命周期事件通过有界 broadcast 通道发出。慢速订阅者会落后并丢失事件,而不会拖住流的打开或关闭。

这条规则针对的是计数和事件,而不是所有共享结构。打开和关闭与正在进行的扫描共用注册表,这会拖慢它们(见注册表,以及为什么选 scc);此外,open 为用户的流获取 pacers 的读锁,而 set_speed_limits 在遍历每个节流器(pacer)期间持有它的写锁。

Supervisor::tracker() 不经过 actor 就交出 Tracker;crate 文档把它和用量一起列为 “is observable without going through the actor” 的部分。Tracker 自己的文档写着 “Nothing here depends on the config”。actor 只在 Actor::new 中构建它一次,任何应用(apply)都不会替换它;其中唯一跟随 spec(期望状态)变化的是各用户的限速值。

读取方 使用 位置
etemenanki-app 的 REST API flows、sessions、kill、kill_where、subscribe、stats webclient/src/tracking.rs:GET /v1/connections、DELETE /v1/connections/{flow}、DELETE /v1/connections?outbound=…&inbound=…&session=…&user=…、GET /v1/connections/events、GET /v1/traffic
FFI 库 stats、flows、kill ffi/src/proxy.rs → traffic、connections(按 id 排序,最早的在前)、close_connection
katana subscribe src/manager/node.rs → spawn_audit,它根据 FlowEvent::Opened 记录审计命中

REST API 的几处细节直接沿用 Tracker 的语义。GET /v1/connections 按 id 对流排序。DELETE /v1/connections/{flow} 就是 Tracker::kill:只要该 id 的流仍在注册表中(无论之前是否已被强制关闭)就返回 204,没有时返回 404 和 no live flow <id>,例如流的包装器已经被 drop 之后。DELETE /v1/connections 的选择器是对所给字段的 kill_where(一个字段都没给时选中所有流);其中 user 字段匹配用户的 label,并且只匹配有用户的流;未知的查询键会被拒绝并返回 400(deny_unknown_fields)。响应是 {"closed": n},即 kill_where 返回的数量。

REST API 页面描述这些端点,FFI 页面描述这些调用,katana 审计页面描述审计规则。

supervisor/src/track/mod.rs
const EVENT_CAPACITY: usize = 1024;
#[derive(Clone)]
pub struct Tracker {
inner: Arc<Inner>,
}
struct Inner {
flows: scc::HashMap<FlowId, Arc<FlowEntry>>,
next: AtomicU64,
events: broadcast::Sender<FlowEvent>,
/// Flows that closed, for the sampler to fold their last bytes in.
closed: mpsc::UnboundedSender<Arc<FlowEntry>>,
stats: watch::Receiver<Arc<StatsSnapshot>>,
sessions: Arc<Sessions>,
/// Each user's speed limit, shared by their flows.
pacers: parking_lot::RwLock<HashMap<UserKey, Arc<Pacer>>>,
}
impl Tracker {
pub(crate) fn new(sessions: Arc<Sessions>) -> (Self, sampler::Sampler);
pub fn flows(&self) -> Vec<Arc<FlowEntry>>;
pub(crate) fn for_each_flow(&self, visit: impl FnMut(&Arc<FlowEntry>));
pub fn flow(&self, id: FlowId) -> Option<Arc<FlowEntry>>;
pub fn sessions(&self) -> Vec<SessionStats>;
pub fn kill(&self, id: FlowId) -> bool;
pub fn kill_where(&self, select: impl FnMut(&FlowEntry) -> bool) -> usize;
pub fn subscribe(&self) -> broadcast::Receiver<FlowEvent>;
pub(crate) fn set_speed_limits(&self, limits: &HashMap<UserKey, NonZeroU64>);
pub fn stats(&self) -> watch::Receiver<Arc<StatsSnapshot>>;
pub(crate) fn meter_stream<S>(&self, meta: FlowMeta, stream: S) -> Metered<S>;
pub(crate) fn meter_datagram<D>(
&self,
meta: FlowMeta,
link: D,
charge: Option<Arc<Wire>>,
) -> MeteredDatagram<D>;
fn pacer(&self, key: UserKey) -> Arc<Pacer>;
fn open(&self, meta: FlowMeta, charge: Option<Arc<Wire>>) -> FlowHandle;
}

Tracker 是一个廉价的句柄:克隆它只克隆一个 Arc。该模块重新导出了 Metered、MeteredDatagram、StatsSnapshot、TagStats 和 UserStats,所以前端程序统一以 etemenanki_supervisor::track::… 引用它们。

Inner 的字段 类型 内容
flows scc::HashMap<FlowId, Arc<FlowEntry>> 存活流的注册表。为什么用 scc:见注册表,以及为什么选 scc。
next AtomicU64,从 1 开始 下一个流 id。open 用 fetch_add(1, Relaxed) 取号,所以 id 在 supervisor 的整个生命周期内唯一,并按 open 取号的顺序递增。取 id 和插入条目是两个独立的步骤,所以两个并发的打开操作进入注册表的顺序可能与它们的 id 顺序相反。
events 容量为 EVENT_CAPACITY = 1024 的 broadcast::Sender<FlowEvent> 打开和关闭事件。Tracker 只保存发送端;所有接收端都来自 subscribe。
closed mpsc::UnboundedSender<Arc<FlowEntry>> 关闭队列:每条注销的流都会被发送到这里,让采样器统计它自上一个 tick 以来搬运的字节。采样器持有接收端(graveyard),并在每个 tick 把它清空。
stats watch::Receiver<Arc<StatsSnapshot>> 最近一次发布的快照。采样器持有发送端。
sessions Arc<Sessions> 会话注册表。sessions() 读取它,采样器也通过它访问用量账本。
pacers parking_lot::RwLock<HashMap<UserKey, Arc<Pacer>>> 每个用户 key 一个令牌桶。不限速的节流器在 poll_ready 和 charge 中各只花一次原子读取。
方法 作用 开销
flows() 所有存活条目,克隆进一个按 map 的 len() 预分配的 Vec,顺序不定 遍历一次 map(iter_sync)
flow(id) 存活条目 id,如果它仍在注册表中 一次 read_sync
sessions() Sessions::stats():每个存活会话及其线上字节,按会话 id 排序(用户、principal 与会话) 在会话注册表的锁下遍历一次会话,释放锁之后再逐个读取每个会话的 Wire,因为读取 QUIC 会话的字节要获取其连接的锁,而打开和准入不能等待这把锁
kill(id) 强制关闭流 id。只要 id 在注册表中就返回 true,无论之前是否已被强制关闭;其包装器被 drop 之后返回 false。 一次 read_sync
kill_where(select) 强制关闭 select 选中且尚未被强制关闭的每条存活流,并返回它强制关闭的数量 遍历一次;select 在 iter_sync 内部运行
subscribe() 一个接收端,接收从此刻起发出的所有事件 无
stats() watch 接收端的一个克隆 无
set_speed_limits(limits) 设置每个节流器的速率;见限速 在写锁下遍历一次 pacers
meter_stream、meter_datagram 登记一条流并包装它的 link 一次 insert_sync 和一个事件
for_each_flow(visit) 不收集结果地访问每条存活流;即采样器的扫描 遍历一次
supervisor/src/track/mod.rs
#[derive(Debug, Clone)]
pub(crate) struct FlowMeta {
pub session: Option<SessionId>,
pub inbound: CompactString,
pub principal: Arc<Principal>,
pub source: Option<IpAddr>,
pub destination: Destination,
pub sniffed: Option<CompactString>,
pub outbound: OutboundId,
pub rule: Option<RuleId>,
pub plane_epoch: u64,
}
impl FlowMeta {
pub(crate) fn of(
flow: &Flow,
session: Option<SessionId>,
inbound: &CompactString,
listener_source: Option<IpAddr>,
outbound: OutboundId,
rule: Option<RuleId>,
plane_epoch: u64,
) -> Self;
}

FlowMeta 描述一条流是什么,在流打开时固定下来。of 根据路由后的 Flow 和 connector 的上下文构建它:

字段 来源
session connector 的会话 id。connector 没有会话时为 None:TUN 设备的流,因为设备不准入用户、也没有会话;以及在监听器停止时恰好完成握手的 Hysteria 2 连接,run_hysteria_inbound(supervisor/src/serve.rs)会给它一个没有会话的 connector 和一个已经取消的 token
inbound connector 的 FlowContext 中的入站 tag
principal flow.user.user_data:入站准入的 Principal,带有其用户 key 和 label,或者是 Principal::anonymous()
source flow.source;流没有给出时,取监听器上的对端地址
destination flow.destination;对 UDP 子链路来说,是打开它的那个数据包的目的地
sniffed 嗅探恢复出的域名(如果有)。UDP 子链路总是 None,因为它的 Flow 来自 Flow::toward(protocols/src/flow.rs),后者设置 sniffed: None
outbound 流所在版本的 OutboundId。负载均衡器会先解析为某个成员,所以流记录的是成员,而不是负载均衡器。
rule 匹配的规则在该 plane 的 RouteSpec::rules 中的下标;默认路由以及由 supervisor 自己的 DNS 服务应答的 DNS 查询为 None
plane_epoch 为它做路由的 plane 的 epoch
supervisor/src/track/mod.rs
pub struct FlowEntry {
id: FlowId,
meta: FlowMeta,
started: Instant,
// Written by the task relaying the flow alone: uncontended.
up: AtomicU64,
down: AtomicU64,
// What the sampler has counted so far; only the sampler touches these.
sampled_up: AtomicU64,
sampled_down: AtomicU64,
kill: AtomicBool,
read_waker: AtomicWaker,
write_waker: AtomicWaker,
}
impl FlowEntry {
pub fn id(&self) -> FlowId;
pub fn session(&self) -> Option<SessionId>;
pub fn inbound(&self) -> &str;
pub fn user(&self) -> Option<UserKey>;
pub fn user_label(&self) -> &str;
pub fn source(&self) -> Option<IpAddr>;
pub fn destination(&self) -> &Destination;
pub fn sniffed(&self) -> Option<&str>;
pub fn outbound(&self) -> &OutboundId;
pub fn rule(&self) -> Option<RuleId>;
pub fn plane_epoch(&self) -> u64;
pub fn started(&self) -> Instant;
pub fn up(&self) -> u64;
pub fn down(&self) -> u64;
pub fn kill(&self);
pub fn killed(&self) -> bool;
fn take_sample(&self) -> (u64, u64);
}

每条流一个 FlowEntry,放在 Arc 后面,由注册表、包装器、仍在队列中的事件以及任何持有它的调用方共享。访问方法读取 meta:user() 是 principal 的用户 key(匿名时为 None),user_label() 是它的 label(匿名时为空)。started 是流登记时取的 std::time::Instant,此时拨号已经完成。

字段 写入方 内存序 含义
up 中继任务,在每次成功写入或发送之后 Relaxed 到目前为止发往目的地的有效载荷字节数
down 中继任务,在每次成功读取或接收之后 Relaxed 到目前为止从目的地返回的有效载荷字节数
sampled_up、sampled_down 仅采样器,在 take_sample 中 Relaxed 水位线:up 和 down 中已被采样器统计过的部分
kill kill() Release 存储,Acquire 读取 只设置一次,从不清除
read_waker、write_waker 包装器,在每次 poll 返回 Pending 时 futures::task::AtomicWaker 强制关闭时要唤醒的任务

kill() 先存入 true,再唤醒两个 waker。take_sample() 把每条水位线换成当前计数值并返回差值;它是私有的,只有采样器调用。

FlowEntry 的 Debug 是手写的。它打印 id、meta(FlowMeta 派生的 Debug)、up、down 和 killed,计数器和标志都通过访问方法读取;started、水位线和 waker 不打印。

supervisor/src/track/mod.rs
#[derive(Clone)]
pub enum FlowEvent {
Opened(Arc<FlowEntry>),
Closed(Arc<FlowEntry>),
}

两个变体携带的都是条目本身,而不是其计数器的副本。Opened 中的条目在事件发出后仍会继续计数;Closed 中的条目带着最终计数,因为它的包装器已经不在了。手写的 Debug 只打印变体和流 id。

supervisor/src/track/mod.rs
struct FlowHandle {
entry: Arc<FlowEntry>,
tracker: Arc<Inner>,
/// The session wire the flow's payload also counts toward, if any.
charge: Option<Arc<Wire>>,
/// The user's speed limit, charged with the payload; `None` for a flow
/// no user was admitted for.
pacer: Option<Arc<Pacer>>,
}
impl FlowHandle {
fn check(&self) -> Option<io::Error>;
fn add_up(&self, n: usize);
fn add_down(&self, n: usize);
}
impl Drop for FlowHandle { /* deregister */ }
fn killed() -> io::Error {
io::Error::new(io::ErrorKind::ConnectionAborted, "the flow was killed")
}

每个包装器持有的私有登记凭证。流被强制关闭后,check 返回强制关闭错误。add_up 和 add_down 把 n 加到条目的计数器上;有 charge 的 Wire 时也加到它上面(wire.add(n, 0) 或 wire.add(0, n)),有节流器时也记到节流器上。drop 这个句柄就会注销该流。

supervisor/src/track/metered.rs
pub struct Metered<S> {
inner: S,
flow: FlowHandle,
read_pause: Option<Pin<Box<Sleep>>>,
write_pause: Option<Pin<Box<Sleep>>>,
}
impl<S> Metered<S> {
pub fn entry(&self) -> &FlowEntry;
}
impl<S: AsyncRead + Unpin> AsyncRead for Metered<S> { /* poll_read */ }
impl<S: AsyncWrite + Unpin> AsyncWrite for Metered<S> {
/* poll_write, poll_flush, poll_shutdown */
}
pub struct MeteredDatagram<D> {
inner: D,
flow: FlowHandle,
read_pause: Option<Pin<Box<Sleep>>>,
write_pause: Option<Pin<Box<Sleep>>>,
}
impl<D> MeteredDatagram<D> {
pub fn entry(&self) -> &FlowEntry;
}
impl<D: DatagramLink<Addr = Destination>> DatagramLink for MeteredDatagram<D> {
type Addr = Destination;
/* poll_send_to, poll_recv_from */
}

两个构造函数都是 pub(super):只有 Tracker::meter_stream 和 Tracker::meter_datagram 会构建它们。read_pause 和 write_pause 是节流定时器,每个方向一个,只在用户的令牌桶处于欠额时才分配。

supervisor/src/track/pace.rs
const MAX_WAIT: Duration = Duration::from_secs(1);
const MIN_WAIT: Duration = Duration::from_millis(1);
pub(crate) struct Pacer {
/// The current rate; 0 is unlimited, the fast path that takes no lock.
limited: AtomicU64,
state: Mutex<PaceState>,
}
struct PaceState {
/// 0 while unlimited.
rate: u64,
/// Negative while in debt.
tokens: f64,
last: Instant,
}
impl Pacer {
pub(crate) fn unlimited() -> Self;
pub(crate) fn set_rate(&self, rate: Option<NonZeroU64>);
pub(crate) fn charge(&self, n: usize);
pub(crate) fn poll_ready(
&self,
cx: &mut Context<'_>,
pause: &mut Option<Pin<Box<Sleep>>>,
) -> Poll<()>;
}

一个用户的令牌桶:每秒 rate 字节,突发量为一秒的速率。Mutex 是 parking_lot::Mutex,Instant 和 Sleep 来自 Tokio,所以令牌桶运行在 Tokio 的时钟上,测试可以暂停它。限速一节逐个讲解这些方法。

supervisor/src/track/sampler.rs
pub const DEFAULT_SAMPLE_INTERVAL: Duration = Duration::from_secs(1);
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct TagStats {
pub up: u64,
pub down: u64,
pub up_rate: u64,
pub down_rate: u64,
pub flows: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UserStats {
pub user: UserKey,
pub label: CompactString,
pub stats: TagStats,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct StatsSnapshot {
pub tick: u64,
pub interval: Duration,
pub total: TagStats,
pub inbounds: BTreeMap<CompactString, TagStats>,
pub outbounds: BTreeMap<CompactString, TagStats>,
pub users: Vec<UserStats>,
}
pub(crate) struct Sampler {
tracker: Tracker,
graveyard: mpsc::UnboundedReceiver<Arc<FlowEntry>>,
publish: watch::Sender<Arc<StatsSnapshot>>,
totals: Totals,
}
impl Sampler {
pub(crate) async fn run(self, interval: Duration, stop: CancellationToken);
pub(crate) fn tick(&mut self, interval: Duration) -> StatsSnapshot;
}
字段 含义
TagStats::up、down 发往目的地和返回的有效载荷,自 supervisor 启动以来的累计值,包括已关闭的流
TagStats::up_rate、down_rate 上一个 tick 内每秒的字节数:本 tick 统计的字节数除以 interval
TagStats::flows tick 时存活的流数
StatsSnapshot::tick 到目前为止的 tick 数;第一个 tick 之前通道里的默认快照为 0
StatsSnapshot::interval 速率所覆盖的时间:自上一个 tick 以来实际测得的时间,而不是配置的间隔
StatsSnapshot::total 所有流合计,包括匿名流
StatsSnapshot::inbounds 按入站 tag 分组,包括自 supervisor 启动以来承载过流的每个入站
StatsSnapshot::outbounds 按出站 tag 分组,包括每个承载过流的出站。supervisor 自己的 DNS 服务的 tag 是空字符串 ""(它内部的 OutboundId 的 tag 为空,而任何 spec tag 都不可能为空)。Plane::empty()(supervisor 在第一次应用之前持有的 plane)的 blackhole 目标同样是空 tag,版本为 0。
StatsSnapshot::users 有存活流或在上一个 tick 内有流量的用户,按 UserKey 排序。匿名流计入 total、inbounds 和 outbounds,但从不出现在这里。
UserStats::label 协议识别该用户所用的 label,取自采样器统计到的该用户第一条流,之后不再更新

DEFAULT_SAMPLE_INTERVAL 位于 crate 私有的模块中;前端程序看到的是 SupervisorBuilder::sample_interval 的默认值,文档写作 “one second when not set”。etemenanki-app(app/src/instance.rs)、FFI 库(ffi/src/proxy.rs)和 katana(src/manager/node.rs)都保留默认值。

这些 id 来自 supervisor/src/entity/id.rs。FlowId、SessionId、UserKey 和 RuleId 是包装整数的 Copy newtype,带有 new 和 get(都是 const fn);OutboundId 是由 tag 和版本组成的结构体:

类型 包装 Display 在这里的含义
FlowId u64 flow#N 一条流,在一个 supervisor 的生命周期内唯一
SessionId u64 session#N 流所属的会话
UserKey u64 无 supervisor 为前端程序的用户 id 分配的索引;作为节流器和 UserStats 的键
OutboundId tag: CompactString、version: u64 tag@vN 流所在的出站版本
RuleId u32 无 规则在路由 plane 的 RouteSpec::rules 中的下标

在数据平面上,恰好只有两个地方登记流,而且两者都在每个入站拨号必经的路径上。唯一另外调用 Tracker 登记的是基准测试辅助函数 track::bench::metered(基准测试)。

位置 打开什么 包装器 何时登记 charge
supervisor/src/connector.rs → AppConnector::connect,TCP 流 每次 connect 一条 stream:一条连接唯一的流、每条 mux 子流、每条 Hysteria 2 代理 stream、每条 TUN TCP 流、每个 SOCKS CONNECT Metered<Guarded<OutboundStream>> Target::connect_stream resolve 为 Ok 时 None:会话的线上字节在客户端一侧统计
supervisor/src/topology/outbound/udp_fanout.rs → FanOutLink::poll_opening,UDP 关联 关联的数据包被路由到的每个出站版本各一条子链路 MeteredDatagram<Guarded<OutboundDatagram>> 子链路的 connect_datagram resolve 为 Ok 时 connector 以 charging_datagrams() 构建时(SOCKS UDP ASSOCIATE)为会话的 Wire,否则为 None

UDP 流根本不拨号:connect 立即返回一个 FanOutLink,由 fan-out 逐个数据包地路由(出站、UDP fan-out 与负载均衡器)。它打开的每条子链路都是一条独立的流,登记在 connector 交给它的 FlowScope 之下(Tracker、会话 id 和可选的 charge)。fan-out 丢弃一条子链路时会 drop 它的包装器,从而注销它的流;丢弃的原因可能是子链路失败、被强制关闭,或者在超过 MAX_SUBS = 64 时它是最久未使用的那条。子链路表把最近发送过的子链路放在末尾,所以淘汰时移除下标 0,并把接收游标 next 往回移一位,最小到 0。connect_datagram 失败的子链路从不登记;poll_opening 以 debug 级别记录 udp fan-out: opening <tag>@v<version> failed: <error>。

包装器位于 Guarded 之外,所以每次读、写、发送或接收都先检查 kill 标志,再检查节流器,然后是排空守卫,最后才是 link 本身。

flowchart LR
  conn["AppConnector::connect"]
  tcp["Target::connect_stream"]
  fan["FanOutLink"]
  sub["poll_opening: connect_datagram"]
  ms["meter_stream"]
  md["meter_datagram"]
  open["Tracker::open"]
  map["scc 注册表"]
  ev["broadcast: FlowEvent"]
  wrap["Metered 或 MeteredDatagram"]
  conn -->|"TCP"| tcp
  conn -->|"UDP"| fan
  fan --> sub
  tcp -->|"Ok"| ms
  sub -->|"Ok"| md
  ms --> open
  md --> open
  open -->|"insert_sync"| map
  open -->|"Opened"| ev
  open -->|"FlowHandle"| wrap
  1. principal 有用户 key 时,查找该用户的节流器(pacer(key)):在读锁下命中,否则在写锁下插入一个新的不限速节流器。每个用户的流即使没有限速也会拿到节流器,这样之后设置的限速也能作用到已经打开的流。
  2. 从 next 取下一个 id。
  3. 构建条目:计数器和水位线为 0,kill 为 false,started 为当前时刻。
  4. 用 insert_sync 把它插入注册表。结果被忽略;id 从不重复。
  5. 发送 FlowEvent::Opened。发送错误只意味着没有订阅者,会被忽略。
  6. 返回 FlowHandle,由包装器保存。

TCP 流的 FlowMeta 在拨号开始前就已构建,但流只在拨号成功后才登记。connect 每次调用只加载一次 plane(“Loaded once for this flow: a route change reaches the next flow”),所以 FlowMeta 中的路由、负载均衡器的选择和 plane epoch 都是那次加载时确定的。正在进行的拨号或失败的拨号从不进入注册表,不发送事件,也不计数。

  1. 用 remove_sync 把条目从注册表中移除。
  2. 把条目发送到关闭队列,让采样器统计它自上一个 tick 以来搬运的字节。采样器停止后,发送会静默失败。
  3. 发送携带该条目的 FlowEvent::Closed,此时条目上是最终计数。

包装器在以下情况被 drop:运行时 drop 出站时(发生错误、执行关闭 effect 或连接结束之后),fan-out 丢弃子链路时,或连接的任务被 drop 时。离开注册表没有其他途径:条目的存活时间与其包装器完全相同。

计数器记录的是协议核心与出站之间的有效载荷,而不是线上字节:

包装器方法 计数 时机
Metered::poll_write add_up(n),n 为内层写入接受的字节数 仅在 Ready(Ok(n)) 时
Metered::poll_read add_down,数量为 buf.filled() 的增长 仅在 Ready(Ok(())) 时;流结束时加 0
MeteredDatagram::poll_send_to add_up(n) 仅在 Ready(Ok(n)) 时
MeteredDatagram::poll_recv_from add_down,数量为 buf.filled() 的增长 仅在 Ready(Ok(_)) 时
poll_flush、poll_shutdown 不计数

对代理出站来说,这是客户端一侧的协议加帧和加密之前的有效载荷。所以流的计数器、会话的线上字节和上游自己统计的字节三者各不相同;每用户用量统计对比了前两者。

FlowEntry::kill(或 Tracker::kill、Tracker::kill_where)从中继某条流的任务之外关闭这条流。它分两步:以 Release 内存序把 true 存入 kill,然后唤醒 read_waker 和 write_waker。此后包装器在任一方向上的每次 poll,都会在触及节流器或 link 之前返回 Err,错误类型为 io::ErrorKind::ConnectionAborted,文本为 the flow was killed。poll_flush 和 poll_shutdown 也不例外,所以只在等待 shutdown 的运行时同样会得知流已被强制关闭。

包装器的每次 poll 都经过 supervisor/src/track/metered.rs 中两个辅助函数之一:

  • guarded(flow, waker, cx, op):流已被强制关闭时返回强制关闭错误;否则运行 op。如果 op 返回 Pending,就把 cx 的 waker 注册到 waker 中,然后再检查一次标志,如果此时已被设置,就返回强制关闭错误。
  • paced(flow, waker, pause, cx, op):流已被强制关闭时返回强制关闭错误。如果流有节流器且 Pacer::poll_ready 为 Pending,就注册 waker 并再检查一次标志,做法与上面完全相同。否则继续执行 guarded。

先注册、再做第二次检查,就消除了这个竞态:落在第一次检查与注册之间的强制关闭,要么被第二次检查看到,要么发现 waker 已经注册并唤醒任务。poll_read 和 poll_recv_from 挂在 read_waker 上;poll_write、poll_send_to、poll_flush 和 poll_shutdown 挂在 write_waker 上。因用户欠额而挂起的流同样会被强制关闭唤醒,而不只是在节流定时器触发时才醒来(killing_a_paced_flow_fails_its_pending_write)。

强制关闭在包装器的下一次 poll 时生效,之后发生什么取决于谁持有这个包装器:

sequenceDiagram
  participant C as 调用方
  participant E as FlowEntry
  participant R as 运行时任务
  participant M as Metered
  participant K as 协议核心
  C->>E: kill
  E->>E: 以 Release 设置 kill 标志
  E->>R: 唤醒 read_waker 和 write_waker
  R->>M: poll_read 或 poll_write
  M->>R: Err ConnectionAborted, the flow was killed
  R->>R: drop 该 key 及其 link
  M->>E: FlowHandle 被 drop,流被注销
  R->>K: 该 key 的 Event::OutboundError
  K->>R: Close 该子流,或 ShutdownTransport 与 Finish

运行时像对待任何失败的读写一样处理这个错误:它 drop 该 key 的 slot,从而 drop 包装器,并把 Event::OutboundError 交给协议核心(服务端运行时)。其余由协议核心决定:

流 谁收到错误 结果 debug 日志
HTTP、Shadowsocks、Shadowsocks 2022、Trojan、VLESS 或 VMess 连接唯一的那条 TCP 流 连接的运行时,然后是协议核心 on_outbound_gone:协议核心推入 ShutdownTransport 和 Finish(服务端运行时) http: outbound failed: the flow was killed;shadowsocks: outbound gone: …、shadowsocks-2022: outbound gone: …、trojan: outbound gone: …、vless: outbound gone: …、vmess: outbound gone: …
VLESS、VMess 或 Trojan mux 连接的一条子流(Trojan 客户端通过对 mux 地址发起 CONNECT 进入解复用器,见 protocols/src/mux/mod.rs 中的 is_mux_destination) 载体的运行时,然后是协议核心 Demux::on_outbound_gone:解复用器为该子流发送一个 End 帧,并为它的 key 推入 Close;载体和其他子流继续运行,新的子流仍然可以打开 vless: outbound gone: the flow was killed、vmess: … 或 trojan: …
一条 Hysteria 2 代理 stream 该 stream 的运行时,然后是 Hy2StreamCore 在该 stream 的中继上执行 on_outbound_gone:只为这条 stream 推入 ShutdownTransport 和 Finish;QUIC 连接(即会话)及其其他 stream 继续运行 hysteria2: outbound failed: the flow was killed
一条 TUN TCP 流 该流的运行时,然后是 PassthroughCore on_outbound_gone:为这条 TCP 流推入 ShutdownTransport 和 Finish 无
一个 SOCKS CONNECT SOCKS 驱动器的中继 中继返回该错误,连接关闭 socks connection from Some(<client IP>) ended: the flow was killed,来自 serve_connection
任意入站的一条 UDP 子链路 FanOutLink fan-out 丢弃该子链路;关联及其其他子链路继续运行,之后被路由到那里的数据包会打开一条带新 id 的新流。协议核心不会收到任何事件。 udp fan-out: sending on <tag>@v<version> failed: the flow was killed,或 udp fan-out: sub-link <tag>@v<version> ended: the flow was killed

这些都是 debug 级别的日志。fan-out 丢弃发送失败的那个数据包,并把它报告为已发送,就像一条有损链路那样。

强制关闭不等于关闭会话。要结束整条连接,前端程序用一个 Selector 调用 Supervisor::close,它经过 actor 并取消会话的 token(用户、principal 与会话)。

Tracker::kill(id) 在 read_sync 中调用 kill(),并返回该 id 是否已登记,所以对一条已被强制关闭、但包装器尚未 drop 的流,它会再次返回 true。kill_where 跳过已被强制关闭的流,只统计它自己强制关闭的流,所以用同一个谓词调用两次,第二次返回 0(a_selective_kill_leaves_the_other_flows)。

每次打开和关闭都发送到同一个容量为 EVENT_CAPACITY = 1024 的 tokio::sync::broadcast 通道上。

  • Opened 在条目进入注册表之后发送,所以订阅者在流关闭之前都能用 flow(id) 查到它。
  • Closed 在条目离开注册表并排入采样器的队列之后发送,带有最终计数。
  • 对同一条流,Opened 总是先于 Closed 发送:前者由 open 发送,后者由句柄的 Drop 发送。

subscribe() 返回的接收端只能看到订阅之后发出的事件,所以它可能收到某条流的 Closed,却从未见过它的 Opened。发送从不等待:通道保留最近 1024 个事件,落后更多的接收端在下一次 recv 时得到 RecvError::Lagged(n),n 是它错过的事件数,随后从仍保留的最早事件继续。完全没有接收端时,发送失败,事件立即被丢弃;这个错误被忽略。REST API 把一次落后转换成一个携带 missed 的 lagged 事件。

采样器每个 tick 把各流的计数器汇总成一个 StatsSnapshot。它是唯一通过每条流的水位线把字节标记为已统计的读取方,所以一条流搬运的每个字节都恰好在一个 tick 中被统计。

Tracker::new 与 Tracker 一起创建采样器,Actor::new 把它保存在 Actor::sampler 中,直到 SupervisorBuilder::start 运行它。SupervisorBuilder 的 Default 把 sample_interval 设为 DEFAULT_SAMPLE_INTERVAL,而 Supervisor::start(spec) 就是 Supervisor::builder().start(spec),所以不经 builder 启动的 supervisor 每秒采样一次。

  1. start 在做任何事之前先断言间隔不为零。间隔为零时以 the sample interval must not be zero panic;SupervisorBuilder::sample_interval 的文档写明了这个 panic。

  2. start 应用第一份 spec,然后用 tokio::spawn 启动 sampler.run(sample_interval, background_stop),并把它的 JoinHandle 保存在 Actor::background 中。

  3. 关闭时,Actor::shutdown 按以下顺序执行:

    1. 停止所有监听器;
    2. 取消每个负载均衡器的探测 token;
    3. 关闭承载所有 accept 循环和连接的 task tracker,并在宽限期内等待它;
    4. 取消根 token,然后再次等待 task tracker;
    5. 清空监听器;
    6. 取消 background_stop,并 await 所有后台任务,其中包括采样器。

    因此,采样器的最后一个 tick 出现在 task tracker 上的所有 accept 循环和连接都已结束之后。

Supervisor::shutdown(grace) 通过 actor 执行这一流程。如果所有 Supervisor 句柄都在没有调用 shutdown 的情况下被 drop,actor 的命令循环结束,并自行调用 shutdown(Duration::ZERO),以同样的方式停止采样器。check() 也会构建一个 actor,因而也有一个采样器,但从不运行它。

supervisor/src/track/sampler.rs
let mut ticker = tokio::time::interval_at(Instant::now() + interval, interval);
ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
  • 第一个 tick 在启动后一个间隔才到来,而不是立即到来(the_sampler_publishes_once_per_interval_and_once_when_stopped 在恰好 1 秒时看到 tick 1,在 2 秒时看到 tick 2)。
  • MissedTickBehavior::Delay:因运行时繁忙而推迟的 tick 不会被集中补上。下一个 tick 的速率改为覆盖更长的间隔,因为 interval 是实测的(now - last,基于 Tokio 的时钟),而不是假定的。
  • 每次迭代在一个 select! 中同时等待 stop.cancelled() 和 ticker.tick(),无论哪个先完成都执行一次 tick。如果是停止信号先到,它发布这最后一个 tick 后返回。最后一份快照的速率覆盖自上一个 tick 以来的时间。
  • tick 本身在 tokio::task::spawn_blocking 中运行:采样器被移入闭包,再连同快照一起交还。一次 tick 要扫描所有存活流,并调和每个活跃用户的所有存活会话,而读取 Hysteria 2 会话的字节需要获取其 QUIC 连接的锁,所以这些工作都不在异步 worker 上运行。tick 内部的 panic 会通过 resume_unwind 在采样器的任务中重新抛出。
  • 快照通过 watch::Sender::send_replace(Arc::new(snapshot)) 发布,无论有没有人在监听都会成功。
flowchart TB
  rec["Ledger::reconcile"]
  reset["tick += 1,重置本 tick 计数"]
  scan["for_each_flow:take_sample,计为存活"]
  drain["graveyard.try_recv 直到为空:take_sample,计为已关闭"]
  snap["快照:按实测间隔计算速率"]
  pub["在 watch 通道上发布"]
  rec --> reset --> scan --> drain --> snap --> pub
  1. 调和用量账本。 self.tracker.inner.sessions.ledger().reconcile() 把每个存活会话新增的线上字节移入其用户的待取总量,这样取用量时只需取走总量(每用户用量统计)。
  2. 开始本次 tick。 tick 加一,每个分组本 tick 的字节数和存活流数归零。累计值保留。
  3. 扫描存活流。 对每条已登记的流,take_sample 返回自其水位线以来的字节并把水位线上移;count 把这些字节加到总计、流的入站分组、出站分组(按 tag,不按版本),以及流有用户时该用户的分组上,并在每个分组中把这条流计为存活。
  4. 清空关闭队列。 对自上次清空以来注销的每条流,take_sample 返回它自上次采样以来搬运的字节,count 把这些字节加上去,但不把这条流计为存活。
  5. 生成快照。 每个分组变成一个 TagStats。其速率来自 per_second(bytes, interval):它以 u128 计算 bytes * 1_000_000_000 / interval_nanos,结果放不进 u64 时给出 u64::MAX(unwrap_or(u64::MAX)),间隔为零时给出 0。用户只保留在本 tick 有存活流或有字节的那些,并按 key 排序。

第 3 步和第 4 步的顺序,加上水位线,保证了计数的精确。在扫描期间关闭的流,要么已经被访问过,那么它的队列条目只给出此后搬运的字节;要么还没被访问,那么它所有未统计的字节都来自队列条目。条目在清空之后才进入队列的流,由下一个 tick 统计。无论哪种情况,水位线的交换都把每个字节恰好交给一个 tick(the_sampler_counts_every_byte_once_by_inbound_outbound_and_user)。

Totals 保存 tick 计数器,以及一个总计 Group、每个入站 tag 一个、每个出站 tag 一个和每个用户 key 一个(附带 label):

supervisor/src/track/sampler.rs
struct Group {
up: u64,
down: u64,
tick_up: u64,
tick_down: u64,
flows: usize,
}

分组在首次使用时创建,所以 inbounds 和 outbounds 也会列出承载过流、现在空闲的 tag,速率为零。查找 tag 时先用 &str,所以已知的 tag 不产生分配。用户的 label 在创建该用户的分组时保存,取自采样器统计到的该用户第一条流,之后不再更新。

StatsSnapshot 按流统计有效载荷;用量按会话统计线上字节。两者无法互相推导;每用户用量统计中的表格把它们并列对比。

用户的限速是 UserSpec::speed_limit: Option<NonZeroU64>,单位是每秒有效载荷字节数,由该用户的所有流在两个方向上共享,TCP 和 UDP 都算在内。None 表示不限速。修改限速不会改变用户的准入方式,所以用户的会话保持打开(a_speed_limit_change_keeps_the_users_sessions)。

限速值只在一次应用或一次用户编辑提交(commit)时才会到达 Tracker:

  1. Actor::commit(每次应用,包括 start 中的第一次)和 Actor::edit_users(在 set_users、upsert_user 或 remove_user 之后)在采用新的用户 key 之后立即调用 publish_speed_limits(&spec)。
  2. publish_speed_limits 遍历新 spec 中每个用户集的每个用户。它跳过没有 speed_limit 或没有用户 key 的用户;出现在多个用户集中的用户取其中最小的限速。
  3. Tracker::set_speed_limits(&limits) 获取 pacers 的写锁,把每个已有节流器设为其用户的限速,用户不在 limits 中时设为不限速;然后为每个尚无节流器的受限用户创建一个节流器,并设为其速率。

之所以只在提交时发布,是因为被拒绝的应用或被拒绝的用户编辑不能改变正在生效的限速。每条存活流都已持有其用户的 Arc<Pacer>,所以新速率在它下一次 poll 时就会生效,无需重新打开任何东西。crate 文档的说法是:用户的限速 “takes effect on live flows at their next read or write, without ending their sessions”。

限速值由各个前端程序设置。etemenanki-app 不设置任何限速:app/src/lower.rs 在把每个用户降为 spec 时都使用 speed_limit: None。katana 根据节点速率和用户速率设置每个用户的限速,取两者中非零的较小值(src/lower/mod.rs → determine_rate);katana 的 lowering(降为 spec)页面和限速指南说明了这些速率的来源。

PaceState::refill(now) 按当前速率为自 last 以来的时间补充令牌,上限为一秒的速率,并把 last 移到 now。不限速的状态只移动 last。经过的时间是 now.saturating_duration_since(last),所以晚于 now 的 last 不补充任何令牌。

速率未变时,Pacer::set_rate(rate) 什么都不做,所以每次应用都发布相同的限速值不会改动任何令牌桶。否则,它先按旧速率补充,再设置余额:

变化 新余额 原因
从受限变为不限速 0 免除所有欠额
从不限速变为受限 new(满桶) 新受限的用户从满额的突发量开始
从一个限速变为另一个限速 min(balance, new) 以旧速率攒下的余额保留,但不超过新的突发量;欠额保留

然后它以 Release 内存序把新速率存入 limited。

Pacer::charge(n) 接收一次读或写已经搬运的 n 字节。limited 为 0 时,它在一次 Acquire 读取后返回,不加锁。否则它加锁,并在锁内再次检查 rate,为 0 时返回(在读取和加锁之间,可能已经有一次改为不限速的 set_rate 运行过);然后补充令牌并减去 n,即使减到零以下也照减。

Pacer::poll_ready(cx, pause) 决定下一次读或写能否开始:

  1. limited 为 0 时,清空 pause 并返回 Ready。
  2. 否则加锁并补充令牌;速率为 0 或余额不为负时返回 Ready(并清空 pause)。
  3. 否则把等待时间算作 -tokens / rate 秒,限制在 MIN_WAIT = 1 ms 到 MAX_WAIT = 1 s 之间,然后释放锁。把 pause 的 Sleep 设到截止时间 now + wait:第一次用 sleep_until 分配一个,每一轮都对它调用 reset(deadline)。然后 poll 它。如果仍在等待,返回 Pending;如果已经到期,回到第 2 步。

MAX_WAIT 限制了每次休眠的时长,让提高或取消的限速在一秒内生效(a_raised_or_removed_limit_is_felt_within_a_second);处于欠额的流至少每秒醒来一次重新检查。MIN_WAIT 避免极少量的欠额让定时器空转。

flowchart TB
  poll["poll_read、poll_write、poll_send_to 或 poll_recv_from"]
  k1{"已被强制关闭?"}
  p{"流有节流器?"}
  ready{"Pacer::poll_ready"}
  park["注册 kill waker,再查一次 kill,Pending"]
  op["poll Guarded,再 poll link"]
  moved{"Ready Ok?"}
  add["add_up 或 add_down:计数器、wire、Pacer::charge"]
  err["Err: the flow was killed"]
  out["返回结果"]
  poll --> k1
  k1 -->|是| err
  k1 -->|否| p
  p -->|否| op
  p -->|是| ready
  ready -->|"Pending,处于欠额"| park
  ready -->|Ready| op
  op --> moved
  moved -->|是| add --> out
  moved -->|否| out

等待发生在操作之前,扣费发生在操作之后。所以被扣住的读会把字节留在出站里(TCP 目的地的发送在其窗口中积压,UDP link 的数据包留在其 socket 中),被扣住的写返回 Pending,使运行时的转发保持排队,从而停止读取客户端(服务端运行时)。节流定时器和 kill waker 注册的是同一个任务 waker,所以两者中任何一个都能唤醒挂起的 poll。

运行时和 fan-out 都不按固定大小搬运字节:一次读取的大小取决于 staging 的剩余空间,一次写入则是协议核心转发的任意区间,两者都受各协议缓冲区大小的限制。如果限速器要等令牌桶攒够 n 个令牌才搬运 n 字节,就无法在不超出速率的前提下放行一个大于一秒突发量的数据块,而且会让不同协议得到不同的速度。

节流器的做法是:

  1. 只要余额不为负,就允许一次传输开始;
  2. 在传输完成后按其完整大小扣费,即使扣到零以下;
  3. 在补充的令牌还清欠额之前,不允许该用户的任何流开始下一次传输。

无论一次搬完还是分多次搬,每个字节都按速率付费。唯一的余量是突发量(最多一秒的速率,在空闲时攒下)以及余额跌破零时已经在进行中的传输;它们由之后的传输来偿还。下面是来自 a_raised_or_removed_limit_is_felt_within_a_second 的算例,速率为每秒 1000 字节:

步骤 之前余额 搬运 之后余额 下一次传输
设置限速(从不限速变为 1000 B/s) 0 无 1000 立即
一次写入 11 000 字节 1000 全部 11 000 字节 −10 000 10 秒后,期间以每次最多 1 秒的休眠反复检查
0.1 秒后取消限速 −9900 无 0 当前休眠结束时,1 秒以内

保持限速不变时,同样的算法得出 a_limited_users_flow_moves_at_its_limit 的结果:以 100 000 B/s 传输 300 000 字节再加 1 字节,耗时在 2 到 3 秒之间,即一秒的突发量加上 200 000 字节的欠额。

add_up 和 add_down 扣的是同一个节流器,而且一个用户的每条流都持有同一个 Arc<Pacer>,不论其入站、会话或方向。在 100 000 B/s 的限速下,同一用户的两条流一条上传 150 000 字节、一条下载 150 000 字节,它们共享一份突发量和一份欠额,所以再多 1 字节要等 2 秒后才能搬运(a_users_flows_share_one_limit)。不限速用户的流走快速路径,在 poll_ready 和 charge 中各对 limited 做一次原子读取;匿名流(开放的入站、共享凭据、TUN 设备)根本没有节流器(an_unlimited_user_and_an_anonymous_flow_are_not_paced)。受限用户的流每次读或写要获取两次节流器的互斥锁,poll_ready 一次,charge 一次。

每条流打开和关闭时,都会从各个 worker 线程写入注册表;每个采样器 tick、flows() 和 kill_where 都会完整扫描它。它必须在扫描进行时仍让打开和关闭保持快速。在 supervisor/Cargo.toml 中,同一条注释同时说明了 scc 和 futures:“The live-flow registry, and the waker a flow is killed through.” Cargo.lock 把 scc 解析为 3.8.8。

引入 Tracker 的那次 Git 提交记录了选型时的对比:在 32 核机器上,8 个线程在 500 000 个存活条目上反复打开和关闭。

Map 空闲时的打开与关闭吞吐量 与全量扫描并行时
scc::HashMap 105 M ops/s 62 M ops/s(一次扫描耗时 5.4 ms)
dashmap 86 M ops/s 10 M ops/s
papaya 11 M ops/s 未记录
Mutex<HashMap> 5 M ops/s 15 k ops/s

这项对比不属于仓库中的基准测试。

supervisor/benches/tracking.rs 是一个 criterion 基准测试。supervisor/Cargo.toml 把它声明为 [[bench]] name = "tracking",并设置 harness = false;criterion 是 dev-dependency。运行方式:

终端窗口
cargo bench -p etemenanki-supervisor --bench tracking

它有三个组。两个用量组(take_usage, 50k active users 和 reconcile (sampler tick), 50k active users)在每用户用量统计中介绍。第三个组测量一条被计量的流在每字节上增加的开销:

  • 组 relay 256 MiB in 16 KiB reads,使用 Throughput::Bytes(256 << 20) 和 sample_size(20),运行在 current-thread Tokio 运行时上。
  • untracked 以每次 16 KiB 的读取读空 tokio::io::repeat(0).take(256 MiB),就像运行时读取出站那样。
  • metered 读空同一个 reader,但它经过 track::bench::metered 包装。这是一个 #[doc(hidden)] 的公开辅助函数(“What the benches need of the internals; not part of the API”),它把 reader 登记为一个独立 Tracker 中的一条流。该辅助函数构建一个全新的 Sessions(新的 CancellationToken、新的 TaskTracker 和默认账本),调用 Tracker::new 并丢弃采样器,然后用固定的 FlowMeta 通过 meter_stream 计量这个 stream:没有会话,入站 bench,Principal::anonymous(),没有来源,TCP 目的地 bench:443,未嗅探,出站 bench@v1,没有规则,plane epoch 为 0。

被计量的流没有节流器,也没有订阅者,所以差别只在于每次读取时的 kill 检查和计数器累加。加入这个基准测试时,测得 untracked 为 153.3 GiB/s,metered 为 153.1 GiB/s(release 构建,32 核):每次 poll 的额外工作淹没在一次 16 KiB 复制的开销里。

不变量 由谁保证 由哪些测试固定
流恰好在其包装器存活期间位于注册表中 open 在包装器存在之前插入;FlowHandle::drop 移除 a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped、killing_a_flow_fails_its_pending_read
只有已打开的出站才是流 meter_stream 在 connect_stream resolve 为 Ok 之后运行;meter_datagram 在子链路打开之后运行 由构造保证
流统计的恰好是实际搬运的有效载荷 只在 Ready(Ok) 之后按内层结果的大小调用 add_up/add_down a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped(上行 5、下行 11)、mux_sub_flows_and_their_session_count_exactly_what_moved、a_udp_association_counts_each_sub_link_and_charges_its_session
每条 mux 子流和每条 UDP 子链路都是一条独立的流,归属于连接的会话 运行时按 key 调用 connect;fan-out 用其 FlowScope 按子链路计量 mux_sub_flows_and_their_session_count_exactly_what_moved、a_udp_association_counts_each_sub_link_and_charges_its_session(两个出站、两条流、一个会话)
Opened 先于 Closed,且 Closed 带有最终计数 open 在插入后发送 Opened;drop 在移除后发送 Closed a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped
打开和关闭从不等待订阅者 有界 broadcast;忽略发送错误 无
强制关闭让流在任一方向上的下一次 poll 失败,并唤醒已挂起的 poll 注册 waker 前后各检查一次标志;kill 唤醒两个 waker killing_a_flow_fails_its_pending_read、killing_a_paced_flow_fails_its_pending_write
一次强制关闭只结束一条流 其他流有各自的标志;协议核心只关闭失败的 key a_selective_kill_leaves_the_other_flows、killing_one_mux_sub_flow_leaves_its_siblings_and_the_carrier
kill_where 对每条流只计一次 跳过 killed() 的条目 a_selective_kill_leaves_the_other_flows
每个字节恰好在一个 tick 中被统计,无论流存活还是已关闭 只由采样器写的水位线;先扫描再清空 the_sampler_counts_every_byte_once_by_inbound_outbound_and_user
速率是本 tick 的字节数除以传给 tick 的间隔 per_second 除以 interval the_sampler_counts_every_byte_once_by_inbound_outbound_and_user(0.5 秒 50 字节即 100 B/s)
该间隔是自上一个 tick 以来实测的时间 run 传入基于 Tokio 时钟的 now - last 无;测试以显式间隔调用 tick
本 tick 既无存活流也无流量的用户不出现在 users 中 Group::active 过滤 the_sampler_counts_every_byte_once_by_inbound_outbound_and_user
每个间隔发布一次,第一次在一个间隔之后,停止时再发布一次 interval_at(now + interval, …);stop 在 select! 中胜出时再执行一次 tick the_sampler_publishes_once_per_interval_and_once_when_stopped
最后一份快照包含每条流的最后字节 采样器只在所有连接结束之后才停止 mux_sub_flows_and_their_session_count_exactly_what_moved 在流存活期间检查统计;关闭顺序没有专门的测试
限速是字节数的属性,而不是数据块大小的属性 搬运后扣费,可以扣成欠额;下一次之前等待 a_limited_users_flow_moves_at_its_limit
每个用户一个令牌桶,覆盖两个方向和所有流 每个用户 key 一个 Arc<Pacer>;两个 add_* 都向它扣费 a_users_flows_share_one_limit
不限速的流和匿名流从不被延迟 limited == 0 快速路径;没有用户 key 就没有节流器 an_unlimited_user_and_an_anonymous_flow_are_not_paced
提高或取消的限速在 1 秒内生效 MAX_WAIT 限制每次节流休眠 a_raised_or_removed_limit_is_felt_within_a_second
修改限速保留用户的会话 限速不属于准入的一部分 a_speed_limit_change_keeps_the_users_sessions(supervisor/tests/hot_swap.rs)
被拒绝的应用或用户编辑不会改动生效中的限速 publish_speed_limits 只在提交阶段运行 无
重新发布未变的限速不会改动令牌桶 速率相同时 set_rate 提前返回 无
出现在多个用户集中的用户得到最小的限速 publish_speed_limits 中的 and_modify(min) 无
情形 结果
一条 TCP 流拨号失败 不登记任何流;运行时把 ConnectFailed 交给协议核心
一条 UDP 子链路打开失败 不登记流,也不计数;poll_opening 以 debug 级别记录 udp fan-out: opening <tag>@v<version> failed: <error>(出站、UDP fan-out 与负载均衡器)
会话拒绝该流(Session::admit) connect 在路由之前就失败;不登记任何东西(用户、principal 与会话)
流被强制关闭 之后的每次 poll 都以 ConnectionAborted、the flow was killed 失败;见连接如何处理这个错误
排空策略关闭了流所在的出站版本 包装器内部的 Guarded 以 the outbound this flow was opened on was drained 让流失败(plane);运行时或 fan-out 像处理强制关闭一样处理这个错误,流在其包装器被 drop 时注销
没有订阅者 事件发送失败并被忽略
订阅者落后超过 1024 个事件 它的下一次 recv 返回 RecvError::Lagged(n);它从保留的最早事件继续
采样器已停止 向关闭队列的发送失败并被忽略;stats() 的接收端保留最后一份快照,changed() 报告发送端已不存在
SupervisorBuilder::sample_interval(Duration::ZERO) start 在应用任何东西之前以 the sample interval must not be zero panic
tick 内部 panic 在采样器任务中通过 resume_unwind 重新抛出;任务结束,不再发布快照

与数据平面上的一切一样,取消通过 drop 完成。包装器不 spawn 任何东西:drop 一个包装器会 drop 它的 link、节流定时器和 FlowHandle,后者注销该流。采样器只会通过 background_stop 停止(actor 在关闭流程的最后取消它),或者在 tick panic 时停止(见上表)。

常量或上限 值 定义位置 含义
EVENT_CAPACITY 1024 supervisor/src/track/mod.rs 订阅者在丢失事件之前最多可以落后的事件数
DEFAULT_SAMPLE_INTERVAL 1 秒 supervisor/src/track/sampler.rs 未设置 sample_interval 时的 tick 间隔;所有前端程序都保留它
MAX_WAIT 1 秒 supervisor/src/track/pace.rs 流再次查看令牌桶之前的最长节流休眠
MIN_WAIT 1 ms supervisor/src/track/pace.rs 最短节流休眠
突发量 用户速率的一秒 supervisor/src/track/pace.rs → PaceState::refill 用户在空闲时最多能攒下的量
MAX_SUBS 64 supervisor/src/topology/outbound/udp_fanout.rs 一个 UDP 关联保留的子链路数,也就是流数;最久未发送过的先被淘汰
关闭队列 无界 mpsc,每个 tick 清空 supervisor/src/track/mod.rs 自上一个 tick 以来每关闭一条流就有一个 Arc<FlowEntry>

每条存活流占用一个注册表条目和一个 Arc<FlowEntry>(其中有它的 FlowMeta、四个 AtomicU64、一个 AtomicBool 和两个 AtomicWaker),加上包装器的 FlowHandle(最多再多三个 Arc),以及仅在用户处于欠额时每个方向一个 box 起来的 Sleep。没有任何逐字节的工作;每次读或写,不限速的流会增加两次对 kill 标志的 Acquire 读取(link 返回 Pending 时为三次:paced 检查一次,guarded 检查一次,guarded 注册 waker 后再检查一次),有节流器时再加两次对节流器速率的读取,以及一次 relaxed 原子加法;有 charge 的 Wire 时再多一次。

单元测试位于 supervisor/tests/unit/track.rs,作为 track::tests 编译进 crate。它们在全新的 Sessions 上构建 Tracker,并用手工构造的 FlowMeta 计量 tokio::io::duplex stream。采样器测试以显式间隔直接驱动 Sampler::tick,或者在暂停的时钟上 spawn Sampler::run。运行采样器的测试和节流测试使用 #[tokio::test(start_paused = true)],这需要 Tokio 的 test-util feature;supervisor/Cargo.toml 在其 tokio dev-dependency 上启用了它。

测试 固定的行为
a_metered_stream_counts_its_payload_and_leaves_the_registry_when_dropped 登记时发出 Opened;统计上行 5 字节、下行 11 字节;drop 时流离开注册表;Closed 带有最终计数
killing_a_flow_fails_its_pending_read 强制关闭以 ConnectionAborted 唤醒挂在一个永不应答的 link 上的读;包装器消失后 kill 返回 false
a_selective_kill_leaves_the_other_flows 按出站 tag 的 kill_where 强制关闭三条流中的两条,第二次计数为 0,第三条流仍能写入
the_sampler_counts_every_byte_once_by_inbound_outbound_and_user 三个 tick 中总计、入站、出站和用户各自的累计值、速率和存活数;已关闭流的最后字节只统计一次;空闲用户被略去
the_sampler_publishes_once_per_interval_and_once_when_stopped 在暂停的时钟上 spawn run:tick 出现在 1 秒和 2 秒,停止时出现 tick 3
a_limited_users_flow_moves_at_its_limit 以 100 000 B/s 传输 300 000 字节再加 1 字节,耗时至少 2 秒且少于 3 秒
an_unlimited_user_and_an_anonymous_flow_are_not_paced 用户 1 受限时,另一个用户和一条匿名流无延迟地搬运 300 000 字节
a_users_flows_share_one_limit 同一用户的一次上传和一次下载共享一份突发量和一份欠额
a_raised_or_removed_limit_is_felt_within_a_second 欠额达十秒的流在限速取消后 1 秒内恢复
killing_a_paced_flow_fails_its_pending_write 强制关闭在一秒内结束一个因用户欠额而挂起的写

集成测试位于 supervisor/tests/tracking.rs,在回环 socket 上运行一个真实的 supervisor。它们用两个辅助函数等待 supervisor:eventually 每 10 ms 轮询一次探测函数,直到它返回一个值,超过 SOON = 5 秒则以 the condition never held 失败;sessions_gone 用它等待 tracker().sessions() 变空。

测试 固定的行为
mux_sub_flows_and_their_session_count_exactly_what_moved 三条 VLESS mux 子流是同一个会话下的三条流,各自带有其入站、label、出站 direct、无规则,并恰好统计各自的有效载荷;快照中的用户行达到有效载荷之和(每 50 ms 采样一次);连接结束后这些流消失
killing_one_mux_sub_flow_leaves_its_siblings_and_the_carrier 被强制关闭的子流结束;它的兄弟子流仍能回显;新的子流可以打开;会话数保持为 1
a_udp_association_counts_each_sub_link_and_charges_its_session 一个被路由到两个出站的 SOCKS 关联是两条 UDP 流,每个出站一条,每个方向分别为 10 和 5 字节,规则分别为 None 和规则 0,同属一个会话
a_hysteria2_session_bills_its_quic_connections_bytes Hysteria 2 stream 的流在每个方向上恰好统计其 20 000 字节,而会话统计得更多

同一文件中的 usage-sink 测试和 a_hysteria2_user_set_admits_by_password 分别在每用户用量统计和用户、principal 与会话中介绍。supervisor/tests/hot_swap.rs → a_speed_limit_change_keeps_the_users_sessions 在一个存活的 mux 会话上通过 set_users 设置限速,并检查会话 id 不变。

webclient/tests/tracking.rs 通过 REST 路由器端到端地测试 Tracker,所用 supervisor 有一个 SOCKS 入站和一个 direct 出站,请求用 oneshot 发送:

测试 固定的行为
connections_lists_live_sessions_and_flows 一个会话和一条流;会话的线上字节包含 SOCKS 问候和请求,而流在每个方向统计 5 字节的有效载荷;流记录了它的会话、tcp、入站、direct 和无规则
deleting_a_flow_closes_it_alone DELETE /v1/connections/{flow} 返回 204 并关闭那条连接;另一条仍能回显;对同一 id 再次删除返回 404
deleting_by_selector_closes_the_matching_flows 不匹配任何流的选择器关闭 0 条;未知查询键返回 400;outbound=direct&inbound=socks-in 关闭 2 条;没有选择器时关闭所有流
traffic_streams_one_event_per_tick 在 100 ms 间隔下,/v1/traffic 每个 tick 发送一个 traffic 事件,tick 编号连续,总计达到发送的 1000 字节
connection_events_stream_opens_and_closes /v1/connections/events 先发送 opened,再为同一 id 发送 closed,每个方向 7 字节

运行方式:cargo test -p etemenanki-supervisor --lib track、cargo test -p etemenanki-supervisor --test tracking 和 cargo test -p etemenanki-webclient --test tracking。

以下行为没有测试:订阅者落后、被强制关闭的 UDP 子链路、只在提交时发布限速、跨用户集取最小限速、未变的限速不改动令牌桶、MIN_WAIT,以及 tick 迟到时采样器的行为。修改这些路径时请补上测试。测试框架见测试。