跳转到内容

监听器与服务循环

源码文件:16 个 · 核对版本 katana v3.0.1 · Etemenanki 596916d
  • katana/src/manager/transport.rs
  • katana/src/manager/proxy.rs
  • katana/src/manager/node.rs
  • katana/src/serve.rs
  • katana/src/inbound.rs
  • katana/src/connector.rs
  • katana/src/meter.rs
  • katana/tests/unit/serve.rs
  • katana/tests/unit/connector.rs
  • katana/tests/unit/e2e.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/config.rs
  • Etemenanki/protocols/src/hysteria/server/datagrams.rs
  • Etemenanki/concepts/src/runtime.rs

本页沿着节点的流量,从已绑定的套接字一直追到每条连接结束。它覆盖一个监听器 generation(一代实例):持有已绑定监听器的 TransportManager、该 generation 所有任务都运行在其下的 Scope、accept 循环及其三个准入信号量、持有节点用户表的 ProxyManager,以及 drive——运行单条连接的运行时并决定其何时结束的循环。

本页面向修改 src/serve.rs、src/manager/transport.rs 或 src/manager/proxy.rs,或者为 stream 节点新增协议的贡献者。流到达 connector(连接器)之后发生的事(准入、路由、审计、计量)见 准入 和 计量 页面;节点如何决定构建或刷新一代实例,见 节点管理器 页面。

服务层只做四件事,其余一切交给内核运行时和 connector:

关注点 负责者 决定什么
绑定与拆除 src/manager/transport.rs → TransportManager 校验、绑定、开始服务;停止、排空并归还端口。
套接字与流的准入 src/serve.rs → accept_loop、serve_socket;src/manager/proxy.rs → ProxyManager::accept_stream 一个节点持有多少套接字、多少套接字处于传输层握手中、多少流处于协议握手中。
单条连接的生命周期 src/serve.rs → serve_stream、drive 构建哪个协议核心、握手截止时间、进度看门狗、退役。
用户表 src/manager/proxy.rs → Tables、ProxyManager::refresh 每条新连接依据什么认证;替换时不触碰监听器。

运行时本身(缓冲区、effect、截止时间计时器)属于 etemenanki-concepts,协议核心及其各阶段的截止时间属于 etemenanki-protocols;见 服务端运行时。katana 只是把它们包起来:选择核心、提供 connector,并在外部观察运行时。

一个节点提供两种监听器形态之一,两者几乎没有共同之处:

Stream 节点 Hysteria 2 节点
套接字 tokio::net::TcpListener std::net::UdpSocket,交给 quinn
传输层 InboundTransport:TCP、TLS、WebSocket 或 gRPC QUIC,位于 Hy2Inbound 内部
由谁 accept katana 的 accept_loop 监听器自己的 Hy2Inbound::run
由谁运行运行时 katana,每个流一个 drive 监听器,每个代理流一个运行时;启用 UDP 时,每条连接的数据报另有一个
用户表 Tables::Stream(ArcSwap<StreamProtocol>) Tables::Hysteria,其认证器被原地替换
已退役用户的连接如何结束 drive 发现租约被取消后返回 该用户的出站拒绝继续传输(Gate),准入拒绝新的流

一代实例拥有的一切都挂在一个 TransportManager 之下。节点管理器同一时刻最多持有一个。

flowchart TB
  TM["TransportManager"]
  SC["Scope: CancellationToken + TaskTracker"]
  PM["ProxyManager(共享 Arc)"]
  AM["握手失败采样器"]
  AL["accept_loop 或 run_hysteria"]
  SS["serve_socket,每个套接字一个"]
  ST["serve_stream,每个流一个"]
  TB["Tables"]
  DP["Dispatcher + Admission"]
  TM --> SC
  TM --> PM
  SC --> AM
  SC --> AL
  SC --> SS
  SC --> ST
  PM --> TB
  PM --> DP

该 generation 的每个任务(包括采样器)都 spawn 在同一个 Scope 中,所以取消这一个 scope 就会停止监听器以及它接受的每条连接。scope 的 token 是一个全新的根 token,不是任何节点级 token 的子 token:必须由所有者调用 TransportManager::shutdown(或 drop 这个 manager)才会触发。

src/manager/transport.rs → TransportManager 就是一代已绑定的监听器:

pub struct TransportManager {
accept: Scope,
proxy: Arc<ProxyManager>,
}
impl TransportManager {
pub async fn start(
node: &NodeInfo,
cert: &CertConfig,
listen_ip: &str,
enable_vless: bool,
sniff: bool,
users: &[UserInfo],
traffic: Arc<NodeTraffic>,
rules: Arc<RuleManager>,
router: Arc<Router<Outbound>>,
node_tag: CompactString,
hysteria: &HysteriaConfig,
) -> io::Result<Self>;
pub fn proxy(&self) -> &Arc<ProxyManager>;
pub async fn shutdown(self);
}
impl Drop for TransportManager {
fn drop(&mut self);
}

对于不需要新监听器的用户集合或节点限速变化,节点管理器通过 proxy() 调用 ProxyManager::refresh。而来自面板的传输层或协议变化、配置文件中对监听地址、证书设置、enable_vless、disable_sniffing、[node.hysteria] 表或路由规则的修改,以及新的出站池,则会替换整个 TransportManager,这会断开该节点上的所有连接。面板无法看到本地修改,所以由节点管理器自己强制生成这一新 generation。

Listener 是尚未开始服务的已绑定套接字。这个枚举是 src/manager/transport.rs 私有的:

enum Listener {
Stream(TcpListener, InboundTransport),
Datagram(std::net::UdpSocket, Hy2Inbound<UserTag>),
}
fn bind_listener(listen: &str, port: u16) -> io::Result<TcpListener>;
fn bind_datagram(listen: &str, port: u16) -> io::Result<std::net::UdpSocket>;

bind_listener 绑定一个 std::net::TcpListener,设为非阻塞,再用 TcpListener::from_std 转换。bind_datagram 绑定一个普通的 std::net::UdpSocket,原样交给 Hysteria 监听器,因为在监听器内部绑定会让第二个套接字与这个套接字相互竞争。

src/serve.rs → Scope 把一个取消 token 和一个 TaskTracker 组合在一起,使拆除时既能停止任务,又能等待任务结束:

#[derive(Clone, Default)]
pub struct Scope {
pub token: CancellationToken,
tasks: TaskTracker,
}
impl Scope {
pub fn new() -> Self;
pub async fn shutdown(&self);
}
pub fn spawn_scoped<F>(scope: &Scope, fut: F)
where
F: Future + Send + 'static,
F::Output: Send;

spawn_scoped 在 tracker 上 spawn 一个任务,用 tokio::select! 让 fut 与 token.cancelled() 竞争。token 触发时 select 完成,fut 被 drop,无论它当时在等待什么。Scope::shutdown 取消 token、关闭 tracker 并等待 TaskTracker::wait,因此只有在该 scope 下 spawn 的每个任务都结束后才会返回。克隆体共享 token 和 tracker,accept_loop、serve_socket 和 accept_stream 正是借此 spawn 到同一个 scope 中。

src/manager/proxy.rs → ProxyManager 持有一个监听器的用户表及其背后的准入。它总是以 Arc<ProxyManager> 形式持有,所有修改都通过内部可变性完成:

pub enum Tables {
Stream(ArcSwap<StreamProtocol>),
Hysteria {
server: Hy2Inbound<UserTag>,
cfg: HysteriaConfig,
},
}
pub struct ProxyManager {
tables: Tables,
sniff: bool,
dispatcher: Arc<Dispatcher>,
traffic: Arc<NodeTraffic>,
preauth: Arc<Semaphore>,
handshake_failures: AtomicU64,
node_tag: CompactString,
}
impl ProxyManager {
pub fn new(
tables: Tables,
sniff: bool,
dispatcher: Arc<Dispatcher>,
traffic: Arc<NodeTraffic>,
node_tag: CompactString,
) -> Arc<Self>;
pub fn dispatcher(&self) -> Arc<Dispatcher>;
pub fn sniff(&self) -> bool;
pub async fn release_listener(&self);
pub fn note_handshake_failure(&self);
pub fn spawn_auth_monitor(self: &Arc<Self>, scope: &Scope);
pub fn accept_stream(
self: &Arc<Self>,
stream: TransportStream,
source: IpAddr,
session: Arc<OwnedSemaphorePermit>,
scope: &Scope,
);
pub fn refresh(&self, node: &NodeInfo, users: &[UserInfo], enable_vless: bool);
pub fn retire_all(&self);
}

Tables 的两个变体是两种不同的机制,而不是同一机制承载两种数据:

  • Tables::Stream 把 stream 节点的 StreamProtocol(src/inbound.rs)放在 ArcSwap 之后。每条连接用 load_full 读取一次,并基于这份快照构建自己的协议核心。刷新时存入一张完整的新表。
  • Tables::Hysteria 持有在用 Hy2Inbound 的一个克隆。监听器就是那个 UDP 套接字,重建它会重新绑定端口,并在每次面板同步时断开所有已连接的客户端。刷新时只通过 Hy2Inbound::set_authenticator 替换其认证器。拆除时,release_listener 也是通过这个克隆拿到在用的 endpoint。保留 cfg 是为了让刷新以与 start 相同的方式构建认证器。

保留 sniff 的原因与 cfg 相同:刷新之后构建的每个核心都必须与最初那些核心以相同方式嗅探。

src/inbound.rs → StreamProtocol 是 stream 节点的核心进行认证时所依据的表:

pub enum StreamProtocol {
Vmess(Arc<AccountValidator<UserTag>>),
Vless(Arc<vless::Validator<UserTag>>),
Trojan(Arc<trojan::Validator<UserTag>>),
ShadowsocksLegacy(Arc<Resolved<UserTag>>),
Shadowsocks2022 {
config: Arc<Ss2022ServerConfig<UserTag>>,
validator: Option<Arc<ss_2022::Validator<UserTag>>>,
},
}

每个变体的数据都是 Arc,因此由它构建核心只是一次引用计数递增,而不是拷贝。这张表如何由面板的节点和用户列表构建,见 入站与出站构建 页面。

src/serve.rs 中的这些函数负责把一个套接字从 accept 带到其最后一个流结束:

pub async fn accept_loop(
tcp: TcpListener,
transport: InboundTransport,
proxy: Arc<ProxyManager>,
scope: Scope,
);
async fn serve_socket(
transport: Arc<InboundTransport>,
sock: TcpStream,
peer: IpAddr,
proxy: Arc<ProxyManager>,
scope: Scope,
session: Arc<OwnedSemaphorePermit>,
stage: OwnedSemaphorePermit,
);
pub(crate) async fn serve_stream(
proxy: Arc<ProxyManager>,
protocol: Arc<StreamProtocol>,
stream: TransportStream,
source: Option<IpAddr>,
preauth: OwnedSemaphorePermit,
);
async fn drive<const BUF: usize, Core, S, Conn>(
stream: S,
core: Core,
connector: Conn,
preauth: OwnedSemaphorePermit,
retired: watch::Receiver<Option<CancellationToken>>,
) -> Result<(), Ended>
where
S: AsyncRead + AsyncWrite + Unpin,
Core: ProxyCoreDecode<Target = Flow<UserTag>, Error = io::Error, TransportAddr = ()>
+ Established,
Conn: Connector<Flow<UserTag>>,
Conn::Datagram: DatagramLink<Addr = Destination>;
pub async fn run_hysteria(
inbound: Hy2Inbound<UserTag>,
socket: std::net::UdpSocket,
dispatcher: Arc<Dispatcher>,
token: CancellationToken,
);

drive 通过一个私有枚举报告连接如何结束,并通过一个私有 trait 向每个核心要求一项能力:

enum Ended {
Handshake(String),
Relay(String),
}
trait Established {
fn is_established(&self) -> bool;
}

established! 宏为 TrojanCore、VlessCore、VMessCore、ShadowsocksCore 和 Ss2022Core(均基于 UserTag)实现 Established,直接转发到核心自身的 is_established。一旦核心的阶段为 Relay 或 Closing,该方法即为 true:请求已解析,流已打开或正在打开。处于嗅探窗口中的核心尚未 established。

connector 通过 src/connector.rs → LeaseSlot 反向联系到连接:

pub type LeaseSlot = watch::Sender<Option<CancellationToken>>;
impl KatanaConnector {
pub fn new(
disp: Arc<Dispatcher>,
source: Option<IpAddr>,
lease: Option<Arc<LeaseSlot>>,
) -> Self;
}

serve_stream 创建一个 watch::channel(None),把 sender 交给它的 KatanaConnector,把 receiver 留给 drive。connector 准入的第一个流用 send_if_modified 把该用户的租约发布到槽中;之后的流不再改动它,因为同一条连接上的所有流都属于同一个用户。租约本身是 Admission 持有的按用户划分的 CancellationToken;见 准入。

sequenceDiagram
  participant C as 客户端
  participant AL as accept_loop
  participant SS as serve_socket
  participant PM as ProxyManager
  participant D as serve_stream 与 drive
  participant K as KatanaConnector
  C->>AL: TCP 连接
  AL->>AL: try_acquire 会话许可,满则拒绝
  AL->>AL: acquire 阶段许可,满则等待
  AL->>SS: spawn_scoped
  SS->>SS: InboundTransport accept(TLS、WebSocket、HTTP/2)
  Note over SS: sink 产出一个流,释放阶段许可
  SS->>PM: accept_stream
  PM->>PM: try_acquire 预认证许可,load_full 读取表
  PM->>D: spawn_scoped
  C->>D: 协议请求
  D->>K: connect(flow):准入,发布租约
  Note over D: 核心 established,释放预认证许可
  K-->>D: 出站已拨号
  D-->>C: 中继,直到运行时结束、租约被取消或看门狗触发

对于 gRPC 传输层,图中间部分会重复:一个套接字对每个 HTTP/2 stream 产出一个流,每个产出的流各自经过 accept_stream。

accept_loop 把传输层包进 Arc,并创建两个与循环同生命周期的信号量,因此新的一代实例从全新的计数开始:

  • sessions,大小为 MAX_LIVE_CONNECTIONS_PER_NODE(65 536);
  • stages,大小为 MAX_TRANSPORT_STAGES_PER_NODE(2048)。

然后它在 tcp.accept() 与 scope.token.cancelled() 之间循环。对每个被接受的套接字:

  1. 占用一个会话名额,否则拒绝。 sessions.try_acquire_owned() 从不等待。如果信号量已空,循环记录一次握手失败,以 debug 级别记录 dropping connection; live connection limit reached,并丢弃该套接字。拒绝是有意为之:在这里等待会让循环完全停止 accept,让内核的 listen backlog 去吸收过载,从而把一个饱和的节点变成一个沉默的节点。该许可被包进 Arc,因为套接字产出的每个流都持有它的一个克隆。

  2. 占用一个传输阶段名额,必要时等待。 stages.acquire_owned() 在与 scope token 竞争的 select! 中 await,因此拆除永远不会被它卡住。等待期间,循环不接受任何新连接;listen backlog 会暂存后续连接,直到某个套接字归还其阶段名额。

  3. 在同一个 scope 下 spawn serve_socket,并带上两个许可。

accept 错误由 should_backoff_accept_error 分类:

pub(crate) fn should_backoff_accept_error(e: &io::Error) -> bool;

ConnectionAborted 和 Interrupted 只涉及单条连接;循环以 debug 级别记录它们(accept error: ...)并立即继续。其他任何类型(例如文件描述符耗尽)都被视为下一次 accept 同样会遇到的状况。循环以 warn 级别记录 accept error, backing off 100ms: ...,并在 backoff_or_cancelled 中休眠 ACCEPT_ERROR_BACKOFF(100 ms),若 scope 被取消则提前返回。循环从不因错误退出;它只在 scope 被取消时结束,结束时 drop 掉监听器。

serve_socket 在套接字上运行 InboundTransport::accept(来自 protocols/src/transports/accept.rs):

pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);

对 TCP、TLS 和 WebSocket,传输层调用一次 sink 后立即返回。对 gRPC,它对每个 HTTP/2 stream 调用一次 sink,并且只在 HTTP/2 连接结束时返回。传输层对 TLS 握手、WebSocket 升级和 HTTP/2 preface 施加自己的 10 s 截止时间(accept.rs 中的 TRANSPORT_HANDSHAKE_TIMEOUT)。

阶段许可放在一个 parking_lot::Mutex<Option<OwnedSemaphorePermit>> 中,在以下两者中先发生的那一刻释放:

  • 第一个流:sink 闭包从 mutex 中取出许可,设置 yielded 标志,并调用 ProxyManager::accept_stream;
  • katana 自己的 TRANSPORT_HANDSHAKE_TIMEOUT(10 s):一个与 tokio::time::sleep 竞争的 select! 取走许可,然后继续 await 同一个 accept future。

第二条路径是为 gRPC 准备的。它的传输层只在客户端打开 stream 时才产出流,而空闲客户端可能永远不打开,所以由计时器而非第一个流来限定这样的套接字被计为处于传输层握手的时长。计时器到期后,连接继续得到服务,并保留其会话名额。

如果 accept 返回错误且从未产出过流,说明套接字的传输层握手失败,serve_socket 会记录一次握手失败。第一个流之后的错误表示连接结束,而不是握手失败,因此只记录日志(inbound transport ended: ...,debug)。

对每个产出的流,accept_stream:

  1. 用 preauth.try_acquire_owned() 占用一个预认证名额。该信号量有 MAX_PREAUTH_STREAMS_PER_NODE(512)个名额。若没有空闲名额,流被立即丢弃,记录一次握手失败,并以 debug 级别记录 node <tag>: dropping a stream; pre-auth limit reached。选择丢弃而非排队,正是为了防止扫描让任务和缓冲区无限增长;
  2. 匹配 Tables::Stream。Hysteria 节点永远不会收到传输层流;另一分支是 debug_assert!(false, ...),在 release 构建中静默返回;
  3. 用 ArcSwap::load_full 加载当前表。流在整个生命周期内保留这份快照,所以之后发生的刷新不会改变这条连接认证所依据的表。同一 gRPC 套接字上后到的流会加载其到达时的当前表;
  4. 在 scope 下 spawn serve_stream。spawn 出的 future 持有会话 Arc 的一个克隆,因此套接字的会话名额会一直保留到它的最后一个流结束。

serve_stream 构建单条连接所需的各个部分,并按表的变体分派:

表变体 核心 缓冲区大小(BUF_SIZE)
StreamProtocol::Vmess VMessCore::new(accounts, now_unix, sniff, source) 32 KiB
StreamProtocol::Vless VlessCore::new(validator, sniff, source) 16 KiB
StreamProtocol::Trojan TrojanCore::new(validator, sniff, source) 16 KiB
StreamProtocol::ShadowsocksLegacy ShadowsocksCore::new(resolved, sniff, source) 20 KiB
StreamProtocol::Shadowsocks2022 Ss2022Core::with_system_clock(config, validator, sniff, source) 32 KiB

BUF_SIZE 是每个核心的关联常量,成为 drive 的 BUF 参数;运行时为每条连接分配三块这个大小的缓冲区(传输层读取、传输层暂存、出站 scratch)。connector 是 KatanaConnector::new(proxy.dispatcher(), source, Some(Arc::new(lease))),其中 source 是监听器看到的对端地址。

drive 返回后,serve_stream 按结果处理:

drive 结果 效果
Ok(()) 不记录日志。
Err(Ended::Handshake(e)) 记录一次握手失败;以 debug 级别记录 inbound handshake failed: <e>。
Err(Ended::Relay(e)) 以 debug 级别记录 connection ended: <e>。

drive 构建 ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress(),把运行时变成一个产出 Traffic 增量的 Stream(concepts/src/runtime.rs)。轮询这个 stream 就是推动连接前进;drop 它就会取消该连接及其打开的一切。

stateDiagram-v2
  [*] --> Handshake
  Handshake --> Established: 核心 is_established
  Handshake --> EndedHandshake: 运行时错误或 EOF
  Handshake --> EndedHandshake: HANDSHAKE_TIMEOUT
  Established --> Established: 增量有字节移动,重置看门狗
  Established --> Finished: 运行时 stream 结束
  Established --> Finished: 租约被取消
  Established --> EndedRelay: 运行时错误
  Established --> EndedRelay: PROGRESS_WATCHDOG
  EndedHandshake --> [*]
  Finished --> [*]
  EndedRelay --> [*]

握手阶段。 一个连上后从不发言的客户端根本不会产生任何运行时事件,而核心只在第一个事件时才启动自己的握手截止时间。所以 drive 从外部把握手当作一个整体来监视:它在 tokio::time::timeout(HANDSHAKE_TIMEOUT, ...) 中轮询 runtime.next(),直到 runtime.core().is_established()。截止时间从流到达 drive 时开始计算,覆盖请求以及任何嗅探窗口。离开这一阶段的三种方式都会变成 Ended::Handshake:

原因 消息
运行时产出错误 该错误的文本
运行时 stream 结束 closed before the handshake completed
超过截止时间 timed out after 10s

核心一旦 established,drive 就 drop 预认证许可。只有握手阶段计入 MAX_PREAUTH_STREAMS_PER_NODE;已建立的连接只计入会话信号量。

中继阶段。 随后 drive 在三个 future 之间 select:

  • runtime.next()。若 moved(&delta),即 transport_rx、transport_tx、outbound_rx 或 outbound_tx 中任一非零,则 Some(Ok(delta)) 重置看门狗。零增量(一次连接、一次关闭、一次截止时间)不算进度。Some(Err(e)) 返回 Ended::Relay。None 表示核心已完成:连接以 Ok(()) 正常结束。
  • until_retired(retired)。它等待槽中出现租约,然后 await lease.cancelled_owned()。它在 changed 之前先调用 borrow_and_update,因此不会错过握手期间发布的租约。在有流被准入之前,这条连接还没有可失去的用户,该 future 永不完成;若 sender 已不存在,它停在 std::future::pending 上。租约被取消时返回 Ok(()):用户已退役,于是连接的运行时、客户端流以及所有出站一起被 drop。
  • 进度看门狗,即一个 tokio::time::sleep(PROGRESS_WATCHDOG),每当增量有字节移动时重置。若它触发,drive 返回 Ended::Relay("nothing moved for 360s")。

看门狗是兜底手段,不是空闲超时。核心自己的 RELAY_IDLE_TIMEOUT(300 s)会先结束一个安静的流,而且是优雅地结束。看门狗负责放弃那些停止移动却不结束的连接,例如客户端已停止读取、而数据仍暂存待发给它的连接,从而释放套接字、出站和三块缓冲区。PROGRESS_WATCHDOG 为 RELAY_IDLE_TIMEOUT.saturating_add(Duration::from_secs(60)),因此永远不会抢在核心的优雅关闭之前。

Hysteria 节点在 katana 中没有 accept 循环。run_hysteria 构建一个 connector 工厂,并把套接字交给监听器:

let make = move |ip| KatanaConnector::new(dispatcher.clone(), Some(ip), None);

工厂的形态由 Hy2Inbound::run(protocols/src/hysteria/server/inbound.rs)固定:

impl<T: Send + Sync + 'static> Hy2Inbound<T> {
pub async fn run<C, F>(
&self,
socket: std::net::UdpSocket,
make_connector: F,
token: CancellationToken,
) -> io::Result<()>
where
F: Fn(IpAddr) -> C + Send + Sync + 'static,
C: Connector<Flow<T>> + Send + 'static,
C::Future: Send,
C::Stream: Send,
C::Datagram: DatagramLink<Addr = Destination> + Send;
}

run 每启动一个运行时,就以客户端地址调用一次工厂:每个代理流一个,每条连接的 UDP 一个。工厂不传入租约槽,因为 katana 并不为每条 Hysteria 连接持有一个能观察租约槽的任务。已退役用户的流之所以结束,一是因为其计量出站拒绝继续传输(租约被取消后,Gate::poll_open 返回带有 the user was retired 的 ConnectionAborted;见 计量),二是因为准入会拒绝他们接下来打开的流。

Hy2Inbound::run 持续服务,直到传给它的 token 被取消或 endpoint 关闭。每个连接任务都位于 run future 所拥有的 JoinSet 中,所以当 scope drop 这个 future 时,所有在用的 QUIC 连接随之而去。监听器自身的准入由 katana 在 build_hysteria(src/inbound.rs)中固定:

设置 值 达到上限时的效果
ListenerConfig::max_connections HY2_MAX_CONNECTIONS = 4096 拒绝传入连接(incoming.refuse()),客户端立即得知。
ServerConfig::circuit_permits Semaphore::new(MAX_LIVE_CONNECTIONS_PER_NODE) = 65 536 代理流以 H3_REQUEST_REJECTED 重置,或者丢弃本会打开新 UDP 会话的数据包。许可由运行该流中继的任务持有,或由 Hy2UdpCore 中的 UDP 会话条目持有,因此覆盖整个 circuit 的生命周期。

run 只在套接字无法成为 quinn endpoint 时返回错误。此时 run_hysteria 以 error 级别记录 hysteria listener failed: <e> 并返回。这一层没有任何东西会重启它,节点管理器也会保留该 TransportManager,所以节点只有在之后某次迫使新 generation 的变化时才会重建监听器,例如来自面板的传输层变化,或 TransportManager 一节列出的某项配置文件修改。

监听器的内部实现,包括其 HTTP/3 认证和流分类,见 Hysteria 2 服务端 页面。

TransportManager::start 在绑定之前运行所有可能失败的步骤,因此错误的配置永远不会留下一个绑定了一半的监听器:

  1. 暂存用户集合:build_user_entries,然后 traffic.prepare(entries)。
  2. build_transport(node, cert):拒绝 REALITY、PROXY protocol、cert.reject_unknown_sni、ACME 模式以及不支持的传输层,并读取证书文件。
  3. build_protocol(node, &valid_users, enable_vless, tag_for):构建 StreamProtocol 表。
  4. bind_listener(listen_ip, node.port)。
  5. Tables::Stream(ArcSwap::from_pointee(protocol)) 和 Listener::Stream(tcp, transport)。

只有在绑定成功之后,start 才会:

  1. 用 traffic.commit(prepared) 提交暂存的用户集合。此时还没有开始服务,所以没有需要退役的用户。NodeTraffic 的生命周期长于各代实例,因此未变化的用户在重建后保留其计数器;见 流量计费;
  2. 构建 Dispatcher(路由器、审计规则、节点 tag,以及一个全新的 Admission)和 ProxyManager;
  3. 创建 Scope,并在其中 spawn 握手失败采样器;
  4. 在该 scope 中 spawn accept_loop 或 run_hysteria。

在绑定及其之前的任何一步出错,start 都会返回,此时什么都没有绑定,也什么都没有提交:traffic.prepare 只是暂存集合,不会触碰注册表。start 不会重试;接下来怎么办由 节点管理器 决定。对节点的第一代实例,节点管理器会在等待一段时间后重试整个启动过程(包括面板读取);等待时间从 1 s 开始翻倍,上限为 60 s 与轮询周期中较短者。重建失败后,节点没有监听器,直到下一个轮询周期构建出新的监听器。

TransportManager::shutdown 是受支持的拆除方式:

  1. self.accept.shutdown().await:取消 scope,并等待其下每个任务都结束。accept 循环、采样器、每个 serve_socket 和每个 serve_stream 都被 drop;对 Hysteria 节点,run future 连同它拥有的每条连接一起被 drop。drop accept 循环会关闭 TCP 监听器。

  2. self.proxy.retire_all():取消 Admission 中每个用户的租约并清空租约表,使仍然存活的计量出站不再移动任何字节。

  3. self.proxy.release_listener().await:在有限时间内归还套接字。对 stream 节点这一步什么都不做。对 Hysteria 节点,它调用 Hy2Inbound::shutdown:以应用错误码 0x100 关闭 endpoint,最多等待 DRAIN_TIMEOUT(3 s)让其空闲,drop 它,然后每隔 RELEASE_POLL(20 ms)尝试绑定同一本地地址,最多持续 RELEASE_TIMEOUT(3 s)。若端口仍被占用,则以 warn 级别记录 hysteria2: <addr> did not come free within 3s 并返回。

之所以需要第 3 步,是因为 quinn 只有在其 driver 任务在没有连接、也没有 endpoint 句柄剩余的情况下被轮询时才会归还套接字,这大约发生在最后一个句柄被 drop 之后一毫秒。节点管理器在 shutdown 返回后立即绑定同一端口;如果不等待,它每次都会输掉这场竞争,节点会一直不可用,直到下一个轮询周期。Hy2Inbound::shutdown 等待的是端口本身,因为那才是必须成立的条件。

impl Drop for TransportManager 取消 accept.token 并调用 proxy.retire_all()。TaskTracker 和 CancellationToken 在 drop 时都不会做任何事,所以若没有这个实现,只要漏掉一次 shutdown 调用,监听器、其 accept 循环以及每条在用连接就会在进程的整个生命周期内一直运行。Drop 不能 await,因此无法等待任务结束,也无法等待 QUIC 端口释放;它是兜底手段,不能代替 shutdown。

shutdown 按值接收 self,所以每次 shutdown 结束时 Drop 也会运行。它的两个动作都是幂等的:token 已经被取消,租约表也已经为空。

不需要新监听器的用户集合或节点限速变化会经过 ProxyManager::refresh,它从不触碰监听器或在用连接:

  1. 暂存新的用户集合并构建替代品:对 stream 节点是新的 StreamProtocol,对 Hysteria 节点是新的 Authenticator。若构建失败,以 error 级别记录 proxy refresh build failed, keeping current: <e> 并返回,运行中的状态保持不变。
  2. 通过 dispatcher.admission.commit(prepared) 提交注册表。这会拒绝已离开用户的新流,并取消已离开用户以及凭据被重新绑定到另一个 uid 的用户的租约,从而经由 drive 或 Gate 结束他们的连接。限速发生变化的用户保留其连接;此后打开的流使用新的速率。
  3. 发布表:stream 节点用 ArcSwap::store,Hysteria 节点用 Hy2Inbound::set_authenticator。

先处理注册表是有意为之。在表被替换之前,已离开的用户仍可依据旧表认证,但准入已经会拒绝他们的流。反过来的顺序则会让新加入的用户依据新表认证成功,然后因为注册表还不认识他们而被拒绝每一个流。完整的论证见 运行时重载 和 准入 页面。

ProxyManager::spawn_auth_monitor 在该 generation 的 scope 下 spawn 一个每秒 tick 一次的任务(立即发生的第一次 tick 会被消耗掉)。每次 tick 把 handshake_failures 交换为零;若计数超过 HANDSHAKE_FAILURE_ALERT_PER_SEC(10),则以 warn 级别记录:

node <tag>: <n> inbound handshake failures in the last 1s (possible handshake scan/DoS or misconfigured clients)

它只做检测,从不阻塞任何东西。计数器由 stream 路径上四处的 note_handshake_failure 累加:

位置 计数内容
accept_loop 因达到在用连接上限而被拒绝的套接字
serve_socket 第一个流之前的传输层失败
ProxyManager::accept_stream 因达到预认证上限而被丢弃的流
serve_stream 任何 Ended::Handshake:协议错误、提前关闭或 HANDSHAKE_TIMEOUT

Hysteria 监听器自行处理握手,并以 debug 级别记录失败的握手(hysteria2: a handshake failed: ...);它们不会进入这个计数器。

不变量 由谁保证 由谁固定
错误的节点配置永远不会只绑定一半。 TransportManager::start 在 bind_listener 或 bind_datagram 之前运行所有构建步骤,并且只在绑定之后提交流量。 构建步骤的拒绝:tests/unit/inbound.rs(例如 reality_rejected、hysteria_requires_a_certificate)。顺序本身没有专门的测试。
被接受的套接字持有其会话名额,直到最后一个流结束。 会话许可是 Arc<OwnedSemaphorePermit>;serve_socket 以及 accept_stream spawn 的每个任务都持有一个克隆。 没有专门的测试。
在会话上限处的过载被拒绝,从不排队。 accept_loop 中的 try_acquire_owned。 没有专门的测试。
套接字持有其传输阶段名额最多 10 s。 阶段许可由第一次 sink 调用或 serve_socket 中的 TRANSPORT_HANDSHAKE_TIMEOUT 分支取走。 没有专门的测试。
每个流依据其到达时的当前表认证。 accept_stream 中的 ArcSwap::load_full;快照被移入 serve_stream。 由 tests/unit/e2e.rs 中的 unchanged_user_survives_user_refresh 间接覆盖,其中一条在用连接在表替换前后持续中继。
沉默的客户端恰好在 HANDSHAKE_TIMEOUT 时结束。 drive 中包住握手循环的 tokio::time::timeout。 tests/unit/serve.rs 中的 a_silent_client_is_dropped_at_the_handshake_deadline。
预认证名额只持有到核心 established 为止。 drive 中紧跟握手循环之后的 drop(preauth)。 由同一测试间接覆盖;没有测试会填满预认证信号量。
不移动任何数据的连接在 PROGRESS_WATCHDOG 之后结束。 看门狗 Sleep,只由 moved 为 true 的增量重置。 tests/unit/serve.rs 中的 a_connection_that_stops_moving_is_dropped。
让用户退役会结束该用户的 stream 连接。 通过 LeaseSlot 发布、由 until_retired 观察的租约。 tests/unit/serve.rs 中的 a_retired_users_connection_ends;tests/unit/connector.rs 中的 the_lease_reaches_the_connection_and_goes_with_the_user;tests/unit/e2e.rs 中的 unchanged_user_survives_user_refresh。
让用户退役会停止其 Hysteria 流,而不影响其他客户端。 准入拒绝新的流;Gate 在租约被取消后拒绝移动字节;监听器不重建。 tests/unit/e2e.rs 中的 a_retired_user_stops_while_the_rest_keep_their_connections。
Hysteria 刷新从不重新绑定套接字。 Tables::Hysteria 加 Hy2Inbound::set_authenticator。 tests/unit/e2e.rs 中的 repeated_user_refreshes_never_disturb_a_live_connection。
拆除只在 spawn 到该 generation scope 中的每个任务都结束后才返回。 Scope::shutdown(取消、关闭、TaskTracker::wait)。 没有测试直接等待 tracker;tests/unit/e2e.rs 中的 route_change_drops_connections 表明重建会断开已打开的连接。
shutdown 返回时,QUIC 端口已释放,或已记录警告。 ProxyManager::release_listener → Hy2Inbound::shutdown。 没有专门的测试。
被 drop 的 TransportManager 永远不会遗留其任务。 impl Drop for TransportManager。 没有专门的测试。
Hysteria 流与 stream 流以相同方式计费。 同一个 KatanaConnector,由 run_hysteria 中的工厂为每个运行时构建。 tests/unit/e2e.rs 中的 a_hysteria_node_relays_and_meters。

每个工作单元因何结束、释放什么:

单元 何时结束 释放什么
accept_loop scope 被取消。accept 错误从不结束它。 TCP 监听器,以及随之释放的两个信号量。
serve_socket 传输层返回:对 TCP、TLS 和 WebSocket 是在那一个流之后,对 gRPC 是在 HTTP/2 连接结束时;或者 scope 被取消。 阶段许可(如果仍持有),以及它持有的会话许可克隆。
serve_stream drive 返回,或 scope 被取消。 运行时(客户端流、出站、缓冲区)、预认证许可(如果仍持有),以及它持有的会话许可克隆。
采样器 scope 被取消。 无其他资源。
run_hysteria scope token 被取消、endpoint 关闭,或 run 失败。 每个 QUIC 连接任务(由 run 的 JoinSet 持有)。Hy2Inbound 中保留的 endpoint 句柄由 release_listener 释放。

取消总是以 future 被 drop 的形式到来:spawn_scoped 让每个任务与 scope token 竞争,所以任何任务都不需要自己检查取消。正在等待任何东西(传输层握手、拨号、阶段许可、看门狗)的任务会在该处被 drop,其许可由析构函数释放。

常量 定义位置 值 作用对象
MAX_LIVE_CONNECTIONS_PER_NODE src/serve.rs 65 536 每一代 stream 监听器同时存活的已接受套接字;同时也是 Hysteria 监听器 circuit_permits 的大小。
MAX_TRANSPORT_STAGES_PER_NODE src/serve.rs 2048 同时处于传输层握手中的套接字。超过后,accept 循环等待。
TRANSPORT_HANDSHAKE_TIMEOUT src/serve.rs 10 s 套接字持有其传输阶段名额的时长。
TRANSPORT_HANDSHAKE_TIMEOUT protocols/src/transports/accept.rs 10 s 传输层自身对 TLS、WebSocket 升级和 HTTP/2 preface 的截止时间。
ACCEPT_ERROR_BACKOFF src/serve.rs 100 ms ConnectionAborted 或 Interrupted 以外的 accept 错误之后的暂停。
MAX_PREAUTH_STREAMS_PER_NODE src/manager/proxy.rs 512 每代实例中同时处于协议握手中的流。超过后,流被丢弃。
HANDSHAKE_TIMEOUT protocols/src/core/mod.rs 10 s 从 drive 开始到核心 established 为止。
RELAY_IDLE_TIMEOUT protocols/src/core/mod.rs 300 s 核心自身对中继中的流的空闲上限。
PROGRESS_WATCHDOG src/serve.rs 360 s(RELAY_IDLE_TIMEOUT + 60 s) 不移动任何字节的已建立连接。
HANDSHAKE_FAILURE_ALERT_PER_SEC src/manager/proxy.rs 10 每 1 s 采样中失败次数超过此值时,采样器发出警告。
HY2_MAX_CONNECTIONS src/inbound.rs 4096 每个 Hysteria 监听器的并发 QUIC 连接数。
DRAIN_TIMEOUT protocols/src/hysteria/server/inbound.rs 3 s 等待已关闭的 endpoint 进入空闲。
RELEASE_TIMEOUT / RELEASE_POLL protocols/src/hysteria/server/inbound.rs 3 s / 20 ms 等待 UDP 端口重新可绑定。

所有信号量都按 generation 创建:会话和阶段信号量在 accept_loop 中,预认证信号量在 ProxyManager::new 中,两个 Hysteria 上限在 build_hysteria 和 Hy2Inbound::run 中。重建会让所有计数从零开始,用户刷新则保留当前计数。

该模块的单元测试位于 tests/unit/serve.rs,作为 src/serve.rs 的 tests 模块编译进去(#[path = "../tests/unit/serve.rs"]),因此可以访问私有的 drive、Ended 和 should_backoff_accept_error。

测试 固定了什么
accept_error_backoff_classification Interrupted 和 ConnectionAborted 不退避;OutOfMemory 和 Other 退避。
a_silent_client_is_dropped_at_the_handshake_deadline 在暂停的时钟下,一个收不到任何字节的 Trojan 核心恰好在 HANDSHAKE_TIMEOUT 之后以 Ended::Handshake 结束。
a_connection_that_stops_moving_is_dropped 一个远端持续灌数据、而客户端从不读取的 passthrough 流,以包含 nothing moved 的 Ended::Relay 结束,且不早于 PROGRESS_WATCHDOG。
a_retired_users_connection_ends 租约有效时连接持续运行;租约被取消后一秒内返回 Ok(())。

该文件中的辅助工具可用于为 drive 编写新测试:

  • impl Established for PassthroughCore<UserTag>,使内核的 PassthroughCore(一个在收到第一批客户端字节时打开固定流的核心)可以被驱动;
  • once(outbound),一个把一条 DuplexStream 出站交给第一个流、拒绝其余流的 connector;
  • flow(),一个发往 192.0.2.1:80 的 TCP 流,用户 tag 为 u、uid 为 1;
  • preauth(),来自一个单名额信号量的许可。

tests/unit/e2e.rs 中的端到端测试针对一个假面板运行整个节点:unchanged_user_survives_user_refresh、route_change_drops_connections、a_hysteria_node_relays_and_meters、a_retired_user_stops_while_the_rest_keep_their_connections 和 repeated_user_refreshes_never_disturb_a_live_connection。katana 中没有测试直接驱动这些上限本身(三个信号量、HY2_MAX_CONNECTIONS)、阶段许可的释放或 QUIC 端口的释放;修改其中任何一项都应附带相应的测试。

在 katana 仓库中运行该模块的测试:

终端窗口
cargo test --locked serve::tests

为 stream 节点新增一个协议,除了在 src/inbound.rs 中编写其构建函数之外,还要在本页涉及的代码中改动三处:

  1. 新增一个 StreamProtocol 变体,以 Arc 持有其表。
  2. 把该核心加入 src/serve.rs 中的 established! 列表,使 drive 能判断其握手何时结束。该核心必须满足 ProxyCoreDecode<Target = Flow<UserTag>, Error = io::Error, TransportAddr = ()>。
  3. 新增一个 serve_stream 分支,调用 drive::<{ NewCore::<UserTag>::BUF_SIZE }, _, _, _>,并使用与其他分支相同的 connector、预认证许可和租约 receiver。

其他都不需要改:准入、看门狗、退役和拆除对 serve_stream 的每个分支一视同仁。