跳转到内容

监听器与服务循环

源码文件:57 个 · 核对版本 Etemenanki 555b7df · katana v4.1.1
  • Etemenanki/supervisor/src/serve.rs
  • Etemenanki/supervisor/src/system/mod.rs
  • Etemenanki/supervisor/src/system/listener.rs
  • Etemenanki/supervisor/src/topology/inbound/mod.rs
  • Etemenanki/supervisor/src/build/inbound.rs
  • Etemenanki/supervisor/src/lib.rs
  • Etemenanki/supervisor/src/supervisor.rs
  • Etemenanki/supervisor/src/connector.rs
  • Etemenanki/supervisor/src/entity/session.rs
  • Etemenanki/supervisor/src/entity/usage.rs
  • Etemenanki/supervisor/src/entity/id.rs
  • Etemenanki/supervisor/src/topology/flow.rs
  • Etemenanki/supervisor/src/topology/spec_plan/inbound.rs
  • Etemenanki/supervisor/src/topology/spec_plan/plan.rs
  • Etemenanki/supervisor/src/build/apply.rs
  • Etemenanki/supervisor/src/build/users.rs
  • Etemenanki/supervisor/src/build/validate.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/transports/accept.rs
  • Etemenanki/protocols/src/transports/tls/config.rs
  • Etemenanki/protocols/src/error.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/http/protocol.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/ss_2022/crypto.rs
  • Etemenanki/protocols/src/ss_aead.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/config.rs
  • Etemenanki/protocols/src/hysteria/server/endpoint.rs
  • Etemenanki/protocols/src/hysteria/obfs.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/tun/device.rs
  • Etemenanki/protocols/src/tun/config.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/app/src/main.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/supervisor/tests/unit/serve.rs
  • Etemenanki/supervisor/tests/unit/session.rs
  • Etemenanki/supervisor/tests/unit/plan.rs
  • Etemenanki/supervisor/tests/unit/validate.rs
  • Etemenanki/supervisor/tests/hot_swap.rs
  • Etemenanki/supervisor/tests/tracking.rs
  • Etemenanki/protocols/tests/pipeline/hysteria.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/tests/unit/trojan/core.rs
  • Etemenanki/app/tests/integration/e2e_unix.rs
  • Etemenanki/app/tests/integration/e2e_hysteria_inbound.rs
  • Etemenanki/app/tests/integration/e2e_tun.rs
  • Etemenanki/app/tests/integration/e2e_route_context.rs
  • katana/src/manager/node.rs

在 supervisor/src/system/listener.rs 和 supervisor/src/serve.rs 中,入站的 spec(期望状态)先变成一个已绑定的 socket,再变成正在运行的连接。监听器在一次应用(apply)的准备(prepare)阶段绑定,所以无法绑定的端口或设备会在任何东西发生变化之前让整个 spec 被拒绝;监听器要等到这次应用提交(commit)时才开始服务。只要它的 BindSpec 还留在 spec 中,它就一直存在:绑定不变而入站有变化时,入站会在同一个 socket 上换一个新的 handler(处理器)。每个被接受的连接都要先通过该监听器的两个信号量的准入,并在传输层握手开始之前登记为一个会话,然后由一个限制握手时长的小看门狗驱动它经过协议核心(core)。Hysteria 2 和 TUN 监听器各自独占一个 UDP socket 或一个设备,并把每个连接的工作交给 protocols crate。

本页面向添加入站协议、修改连接准入方式,或者改动任何运行在监听器 token 或会话 token 之下的代码的贡献者。一次应用要绑定、替换或停止哪些监听器,见规划与应用变更。会话如何绑定到用户、如何被撤销、关闭和列出,见用户、principal(身份主体)与会话。connector(连接器)如何处理每个流,见 plane(数据平面):为每个流选路;会话的字节计数器供给什么,见按用户的用量计费。

关注点 位置 说明
在准备阶段绑定 system/listener.rs → bind、Pending::pair 打开 TCP 或 UDP socket、Unix socket 文件或 TUN 设备,然后构建服务所需的一切。此时还不提供服务。
在提交阶段服务 system/listener.rs → Listener::start spawn 服务一个待启动配对的任务。它不会失败。
监听器身份 topology/inbound/mod.rs → BindSpec 不变的绑定在多次应用之间保留它的 socket 和准入计数。
替换 handler Listener::prepare_swap、Listener::swap 流式监听器:下一个被接受的 socket 得到新 handler。Hysteria 2 监听器在运行中的 endpoint 内部应用它,TUN 监听器则用它重启自己的运行时(见替换 handler)。
用户表 Listener::prepare_users、Listener::store_users;build/inbound.rs → UserTable 用户变化时把一张新表存入运行中的 handler,handler 本身不重建。
接受与准入 serve.rs → run_stream_inbound 每个监听器两个信号量,获取时不等待。
accept 时即为会话 run_stream_inbound → Sessions::open 每个被接受的 socket 在其传输层握手开始之前就已是一个会话。
传输层扇出 serve.rs → Connection::serve_socket InboundTransport::accept 为每个 socket 产出一条字节流;gRPC 则为每条 HTTP/2 stream 产出一条。
协议分派 serve.rs → serve_connection SOCKS 驱动,或在协议核心之上运行的 ProxyServerRuntime。
握手看门狗 serve.rs → drive 在协议核心建立之前,每个运行时步骤都受 HANDSHAKE_TIMEOUT 约束。
Hysteria 2 serve.rs → run_hysteria_inbound 每个 QUIC 连接一个会话;endpoint 原地修改。
TUN serve.rs → run_tun_inbound 每个设备同一时间只有一个运行时,替换时重启;流没有会话。
停止 Listener::stop 停止接受。流式连接和 Hysteria 2 连接继续运行。

这部分代码不解析协议(由 etemenanki-protocols 中的协议核心负责),不选路也不拨号(由它交给每个连接的 AppConnector 负责),不决定绑定、替换或停止哪些监听器(由规划负责),也不会自行结束会话(由谁结束见 token 与任务)。它唯一的计量工作是向每个会话的 Wire 计数器累加字节。

每个前端程序都以同一种方式到达这里:应用一份 spec。etemenanki-app 为自己的配置运行一个 supervisor(监管器),katana 为每个节点运行一个。服务代码本身不公开:serve 是 supervisor/src/lib.rs 中的一个 pub(crate) 模块;system 是公开的,但它唯一的模块 listener 是 pub(crate);build::inbound 也是 pub(crate)。前端程序能看到 topology::inbound 中的类型(BindSpec、TunSource、SuppliedTun 和各 handler 类型),但没有任何公开 API 接受 handler,所以只有本 crate 构建它们(在 build::inbound 中)。

BindSpec:绑定什么,以及监听器的身份

Section titled “BindSpec:绑定什么,以及监听器的身份”
supervisor/src/topology/inbound/mod.rs
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BindSpec {
Tcp { host: String, port: u16 },
Udp { host: String, port: u16 },
Unix(PathBuf),
Tun(TunSource),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TunSource {
Create(DeviceSpec),
Fd(SuppliedTun),
}
#[derive(Debug, Clone)]
pub struct SuppliedTun {
pub fd: Arc<OwnedFd>,
pub mtu: u16,
}

绑定就是监听器的身份。plan 按 bind == 而不是按 tag 把每个期望的入站与一个运行中的入站对应起来,Actor::prepare 也用同样的方式找到运行中的 Listener(listeners.iter().position(|l| l.bind == inbound.bind))。两个绑定的每个字段都相等时才相同。DeviceSpec 比较它的 name、mtu、addresses 和 routes,所以设备的任何变化都是另一个绑定。SuppliedTun 手工实现 PartialEq:持有相同的原始描述符编号和相同的 mtu 时两者相等。代码依赖于描述符编号在打开期间唯一这一点,而两者都保持打开。

校验让绑定成为一个可用的 key。validate 拒绝两个入站共用一个绑定(<bind> is bound by another inbound),validate_inbound 把每种协议与一种绑定配对:SOCKS、HTTP、Trojan、VLESS、VMess、Shadowsocks 和 Shadowsocks 2022 配 Tcp 或 Unix,Hysteria 2 配 Udp,TUN 配 Tun(the protocol cannot be served on <bind>)。同一端口上的 TCP 绑定和 UDP 绑定是不同的绑定。这些规则及其测试见校验与应用错误。

BindSpec 实现了 Display,每条监听器日志和每个绑定错误都用它:

变体 显示为 示例
Tcp {host}:{port} 0.0.0.0:1080
Udp udp {host}:{port} udp 0.0.0.0:443
Unix unix:{path} unix:/run/etemenanki/socks.sock
Tun(Create) tun {name};name 为 None 时为 tun auto tun ete0
Tun(Fd) tun fd {raw fd} tun fd 7

TunSource::Create 是入站自己创建、自己配置地址和路由的设备,在 Linux 上需要 CAP_NET_ADMIN;最后一个描述符关闭时,内核会删除它。TunSource::Fd 是交过来时已经配置好的设备,例如移动端 VPN API 给出的描述符,接口及其地址和路由都归平台所有。这个变体由 FFI crate 构建。监听器服务的是 fd 的副本,停止时关闭这些副本;fd 本身在最后一个持有者(包括 spec)drop 它时关闭。TunSource::mtu() 从任一变体返回 MTU;校验要求至少为 1280(MIN_TUN_MTU)。

supervisor/src/topology/inbound/mod.rs
pub enum InboundHandler {
Stream(Arc<StreamInbound>),
Hysteria2(Hy2Handler),
Tun(TunHandler),
}
pub struct StreamInbound {
pub tag: CompactString,
pub protocol: StreamProtocol,
pub transport: InboundTransport,
pub sniff: bool,
}
pub enum StreamProtocol {
Socks(ArcSwap<SocksInbound<Principal>>),
Http(ArcSwap<HttpServerConfig<Principal>>),
Trojan(ArcSwap<trojan::Validator<Principal>>),
Vless(ArcSwap<vless::Validator<Principal>>),
Vmess(Arc<AccountValidator<Principal>>),
Shadowsocks(ArcSwap<Resolved<Principal>>),
Ss2022(ArcSwap<Ss2022Users>),
}
pub struct Ss2022Users {
pub config: Arc<Ss2022ServerConfig<Principal>>,
pub validator: Option<Arc<ss_2022::Validator<Principal>>>,
}
pub struct Hy2Handler {
pub tag: CompactString,
pub connection: Arc<Hy2ConnectionConfig<Principal>>,
pub quic: quinn::ServerConfig,
pub obfs: Option<Obfs>,
pub max_connections: usize,
}
pub struct TunHandler {
pub tag: CompactString,
pub inbound: TunInbound<Principal>,
}

handler 是监听器提供服务所用的全部东西。它与监听器分开,这样有变化的入站可以在同一个 socket 上换一个新的 handler。

  • StreamInbound 是 accept 循环服务每个 socket 所用的 handler。tag 是路由匹配的依据,也是每个会话登记的归属。transport 是协议之下的那一层;Unix 监听器不带传输层,也从不读取它。sniff 表示协议核心是否从按 IP 寻址的流中读出域名。
  • StreamProtocol 保存构建每个连接的协议核心所需的东西。所有与用户相关的部分都放在 ArcSwap 之后,所以用户变化时只需把一张新表存入运行中的 handler,从不重建它。VMess 在形式上是例外,效果上不是:它保存一个 Arc<AccountValidator>,由它自己替换用户集(AccountValidator::set_users),这样能保留它的重放状态。类型参数是 Principal,即每个协议的用户表在握手成功时交回的按用户载荷(见用户、principal 与会话)。StreamProtocol::name() 返回日志中使用的名字:socks、http、trojan、vless、vmess、shadowsocks 或 shadowsocks-2022。
  • Ss2022Users 是 Shadowsocks 2022 入站的用户表:带有密钥和用户的服务端配置,以及多用户模式下用于身份头(identity header)的查找表。
  • Hy2Handler 是 Hysteria 2 监听器在它唯一的 endpoint 内部应用的内容,因为两个 QUIC endpoint 不能共用一个端口。connection 是每个连接的配置,其中包含用户表(authenticator)和 circuit 预算(circuit_permits)。quic 是回应新握手所用的 TLS 材料。
  • TunHandler 就是 TUN 入站本身,入站是新的或有变化时构建。
supervisor/src/build/inbound.rs
pub(crate) fn build_handler<U: UserId>(
spec: &InboundSpec,
users: &Admissions<U>,
circuits: Option<&Arc<Semaphore>>,
) -> io::Result<InboundHandler>;
pub(crate) enum UserTable {
Socks(SocksInbound<Principal>),
Http(HttpServerConfig<Principal>),
Trojan(trojan::Validator<Principal>),
Vless(vless::Validator<Principal>),
Vmess(Vec<(Uuid, Arc<Principal>)>),
Shadowsocks(Resolved<Principal>),
Ss2022(Ss2022Users),
Hysteria2(Authenticator<Principal>),
None,
}
pub(crate) fn user_table<U: UserId>(spec: &InboundSpec, users: &Admissions<U>) -> io::Result<UserTable>;
impl UserTable {
pub(crate) fn fits(&self, protocol: &StreamProtocol) -> bool;
pub(crate) fn store_into(self, protocol: &StreamProtocol);
}

入站是新的或有变化时,build_handler 构建整个 handler。它先用 user_table 构建用户表,然后:

协议 build_handler 构建什么
Hysteria 2 hy2_handler:伪装响应(Masquerade::default() 或 Masquerade::new(status, body, content_type));circuit_permits,circuits 为 Some 时是运行中监听器的信号量,否则是新建的 Semaphore::new(max_circuits);连接配置(放在 ArcSwap 中的认证器、伪装响应、sniff、由 udp_idle_timeout 得到的 udp、circuit 许可);由 hy2_endpoint::server_config(cert_pem, key_pem) 得到的 quic;从 spec 映射来的 obfs;max_connections。
TUN TunInbound::new(TunConfig { user: Principal::anonymous(), mtu: device.mtu(), udp, udp_idle_timeout, max_flows });除非 spec 开启嗅探,否则再调用 without_sniffing()。
所有流式协议 一个 StreamInbound,用户表放在它的 StreamProtocol 变体中。Unix 绑定上,以及不带传输层的协议(SOCKS、Shadowsocks、Shadowsocks 2022),传输层是 InboundTransport::Tcp。对 HTTP、Trojan、VLESS 和 VMess,build_transport 构建 Tcp、Tls(无 ALPN)、Ws 或 Grpc。TLS 来自 TlsServerConfig::from_pem(cert, key, alpn);叠在 WebSocket 之下时 ALPN 为 http/1.1,叠在 gRPC 之下时为 h2。给定 host 时,InboundTransport::ws(path, host, tls) 要求请求的 Host 与之匹配。

无法解析的证书或私钥在这里失败,成为 ApplyError::Build,显示为 building inbound <tag> failed: <error>。对一个不含任何证书的 Hysteria 2 证书文件,etemenanki-app 的 --test 会显示:

ERROR etemenanki_app: configuration invalid: building inbound hy2-in failed: hysteria2: the certificate file contains no certificates

hy2_endpoint::server_config 构建一个仅支持 TLS 1.3、ALPN 为 h3 的 rustls 配置,以及 Hysteria 2:服务端中描述的 QUIC 传输参数。它的错误有 hysteria2: could not read the certificate: <error>、hysteria2: the certificate file contains no certificates、hysteria2: could not read the private key: <error>、hysteria2: the key file contains no private key、hysteria2: certificate and key do not match: <error>、hysteria2 tls setup failed: <error> 和 hysteria2 quic tls setup failed: <error>。

build_handler 和 user_table 中的其他检查所防范的情况,校验已经先以 ApplyError::Invalid(inbound <tag>: <reason>)拒绝了(plan 在准备阶段之前运行校验):Masquerade::new 检查的伪装响应、normalise_psk 检查的每个 Shadowsocks 2022 用户密钥,以及 inbound <tag>: protocol and bind do not match 和 missing tls 背后的组合。对前两种情况,etemenanki-app 的 --test 打印的是校验给出的文本:

ERROR etemenanki_app: configuration invalid: inbound hy2-in: hysteria2: 233 is the authentication success status and cannot be used for the masquerade
ERROR etemenanki_app: configuration invalid: inbound ss-in: user alice@example.com: shadowsocks-2022: PSK too short (16 < 32)

UserTable 是 handler 中与用户相关的部分,由入站准入的用户构建。每种协议读取哪种凭据、每张表在开放模式或共享模式下保存什么、user_table 填入哪些标签,见用户、principal 与会话。fits 判断一张表是否是某个 StreamProtocol 所存放的那一种,store_into 负责存入:

UserTable 存入者 存入位置
Socks、Http、Trojan、Vless、Shadowsocks、Ss2022 store_into 对相应 StreamProtocol 单元调用 ArcSwap::store
Vmess store_into 对运行中的校验器调用 AccountValidator::set_users
Hysteria2 Listener::store_users Hy2Inbound::set_authenticator
None Listener::store_users 不存:TUN 设备不准入任何用户
supervisor/src/system/listener.rs
pub(crate) enum Bound {
Stream(StreamListener),
Udp(std::net::UdpSocket),
Tun { device: OwnedFd, first: OwnedFd },
}
pub(crate) enum StreamListener {
Tcp(TcpListener),
Unix {
listener: tokio::net::UnixListener,
_file: SocketFile,
},
}
pub(crate) struct SocketFile {
path: PathBuf,
id: (u64, u64),
}
pub(crate) enum AcceptedSocket {
Tcp(tokio::net::TcpStream),
Unix(tokio::net::UnixStream),
}
impl StreamListener {
pub(crate) async fn accept(&self) -> io::Result<(AcceptedSocket, Option<IpAddr>)>;
}
  • StreamListener::accept 对 TCP 返回对端 IP(peer.ip(),端口被丢弃),对 Unix socket 返回 None,因为它没有 IP 对端。
  • drop 一个 StreamListener 就会释放它持有的一切。所以,一次随后被拒绝的应用所绑定的监听器不会留下任何东西;停止的 accept 循环通过 drop 监听器来释放端口。
  • SocketFile 是 Unix 监听器创建的文件系统条目,连同该文件的设备号和 inode。它的 Drop 只在 symlink_metadata(path) 仍然报告同一设备号和 inode 时才删除该路径,所以之后别人放到该路径上的文件不会被动到。删除失败时在 debug 级别记录 could not remove <path>: <error>,除此之外忽略。_file 声明在 listener 之后,所以描述符先关闭,文件后删除。
  • Bound::Tun 携带设备的描述符(由监听器保留)和 first(交给第一个运行时的副本)。
supervisor/src/system/listener.rs
pub(crate) enum Pending {
Stream(StreamListener, Arc<StreamInbound>),
Hysteria2 {
tag: CompactString,
inbound: Hy2Inbound<Principal>,
endpoint: OpenedEndpoint,
circuits: Arc<Semaphore>,
obfs: Option<Obfs>,
},
Tun {
device: OwnedFd,
prepared: PreparedDevice,
handler: TunHandler,
},
}
pub(crate) struct Listener {
pub tag: CompactString,
pub bind: BindSpec,
stop: CancellationToken,
kind: ListenerKind,
}
enum ListenerKind {
Stream(watch::Sender<Arc<StreamInbound>>),
Hysteria2 {
inbound: Hy2Inbound<Principal>,
tag: Arc<ArcSwap<CompactString>>,
circuits: Arc<Semaphore>,
obfs: Option<Obfs>,
},
Tun {
device: OwnedFd,
restart: mpsc::UnboundedSender<(TunHandler, PreparedDevice)>,
},
}
pub(crate) enum Swap {
Stream(Arc<StreamInbound>),
Hysteria2(Hy2Handler),
Tun(TunHandler, PreparedDevice),
}
pub(crate) struct Users(UserTable);
  • Pending 是一个已绑定的句柄与它将要服务的 handler 的配对,服务所需的一切都已构建好。Listener::start 只接受它,所以启动不会失败。
  • Listener 是一个已绑定的句柄加上服务它的任务。tag 是它当前服务的入站的 tag,由 swap 保持最新。bind 是 actor 查找它所用的 key。stop 是 supervisor 根 token 的子 token,所以关停能到达每个监听器。
  • ListenerKind 保存 actor 触达运行中任务所需的东西:流式监听器是其 handler 的 watch 发送端;Hysteria 2 是入站句柄(任务所运行那一个的克隆)、存放当前 tag 的共享单元、circuit 预算和混淆设置;TUN 是设备自己的描述符和运行时重启通道的发送端。
  • Swap 和 Users 分别是已经针对目标监听器检查过的 handler 和用户表。它们在准备阶段生成,在提交阶段消耗。
supervisor/src/serve.rs
#[derive(Clone)]
pub(crate) struct Shared {
pub plane: PlaneCell,
pub sessions: Arc<Sessions>,
pub flows: Tracker,
pub tracker: TaskTracker,
}
impl Shared {
fn connector(&self, tag: CompactString, source: Option<IpAddr>, wire: Wire, stop: &CancellationToken)
-> Option<(AppConnector, Arc<Session>)>;
fn connector_for(&self, tag: CompactString, source: Option<IpAddr>, session: Arc<Session>)
-> AppConnector;
fn spawn_until<F>(&self, token: CancellationToken, fut: F)
where
F: Future<Output = ()> + Send + 'static;
}
struct Connection {
inbound: Arc<StreamInbound>,
peer: Option<IpAddr>,
carrier: Arc<Session>,
stop: CancellationToken,
shared: Shared,
live: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
}
struct WireCounted<S> {
inner: S,
wire: Option<Arc<Wire>>,
}

Shared 是每个 accept 循环与 supervisor 共享的东西,克隆到每个任务中。它的字段:

字段 类型 用途
plane PlaneCell 每个 connector 拨号时经过的路由和出站(plane:为每个流选路)
sessions Arc<Sessions> 会话注册表:accept 时调用 open,Hysteria 2 关停时用 root()
flows Tracker(crate::track) 每个流登记和计量的地方(流跟踪、统计与限速)
tracker tokio_util::task::TaskTracker 每个 accept 循环和连接任务都挂在它上面,以便关停时等待它们

两个 tracker 互不相关:flows 是 supervisor 的流登记表,tracker 是 Tokio 的任务跟踪器。Shared 的方法:

方法 作用
connector 用 Sessions::open 登记一个新会话,并返回用于其流的 connector;监听器的 stop 已取消时返回 None。
connector_for 为已登记的会话构建 connector:AppConnector::new(plane, FlowContext { inbound_tag, source }, flows, Some(session))。
spawn_until 在任务跟踪器上 spawn fut,外面套一个与 token.cancelled() 竞争的 tokio::select!。token 先触发时,future 在它当前所处的位置被 drop。

Connection 是一个被接受的 socket 以及服务它所需的东西:handler、对端 IP、该 socket 自己的会话(carrier)、监听器的 token、Shared,以及两个准入许可(permit)。

WireCounted 包装 SOCKS 驱动的流。poll_read 把一次成功读取填入的字节数计为上行,poll_write 把一次成功写入返回的字节数计为下行,flush 和 shutdown 直接透传。它统计的是什么,见按用户的用量计费。

在准备阶段绑定,在提交阶段服务

Section titled “在准备阶段绑定,在提交阶段服务”
sequenceDiagram
  participant A as Actor(应用)
  participant B as listener bind
  participant P as Pending 配对
  participant L as Listener start
  participant T as 服务任务
  Note over A: 准备阶段,可能失败,外部不可见
  A->>B: 对每个新绑定调用 bind(tag, BindSpec)
  B-->>A: Bound,即 socket、Unix 文件或 TUN 描述符
  A->>P: pair(bound, handler)
  P-->>A: Pending,Hysteria 2 endpoint 已打开或 TUN 设备已准备好
  Note over A: 之后任何失败都会 drop 所有 Pending,从而释放它们
  Note over A: 提交阶段,不会失败
  A->>L: start(bind, pending, shared, 根 token 的子 token)
  L->>T: spawn run_stream_inbound、run_hysteria_inbound 或 run_tun_inbound
  L-->>A: Listener,日志为 inbound tag listening on bind
flowchart TB
  root["根 token(Sessions::root)"]
  stop["监听器的 stop token,根的子 token"]
  sess["会话 token,根的子 token"]
  acc["run_stream_inbound"]
  sock["serve_socket,每个被接受的 socket 一个"]
  conn["serve_connection,每条字节流一个"]
  hy["run_hysteria_inbound"]
  hyconn["QUIC 连接任务(Hy2Inbound::run 中的 JoinSet)"]
  tun["run_tun_inbound"]
  tunflow["流任务(TunInbound::run 中的 JoinSet)"]
  root --> stop
  root --> sess
  stop -. 结束 .-> acc
  stop -. 排空 .-> hy
  stop -. 结束 .-> tun
  acc --> sock
  sock --> conn
  sess -. drop .-> sock
  sess -. drop .-> conn
  sess -. 关闭 .-> hyconn
  hy --> hyconn
  tun --> tunflow

accept 循环、Hysteria 2 和 TUN 监听器任务以及每个流式连接任务,都运行在 supervisor 的 TaskTracker 上。Hysteria 2 监听器的每连接任务和 TUN 设备的每流任务放在 protocols crate 的 run future 所拥有的 JoinSet 里,所以它们随 run 一起结束。监听器的 token 只停止它的 accept 任务。连接运行在它的会话 token 之下,这个 token 是根的子 token,而不是监听器的子 token,所以停止一个监听器之后,它已接受的连接照常运行。服务代码中没有任何地方取消会话的 token。取消它的是:移除策略、Supervisor::close、其入站被移除、关停,或者第一个流出示了一个已被撤销的 principal(见用户、principal 与会话);会话自己的 Drop 也会取消它。TUN 是例外:它的流是运行时的任务,而运行时由监听器的 token 停止。

下面的时序图跟随一个 VLESS over TLS 入站的客户端,从 accept 一直到中继结束。

sequenceDiagram
  participant C as 客户端
  participant L as run_stream_inbound
  participant S as Sessions
  participant K as serve_socket 任务
  participant T as InboundTransport
  participant V as serve_connection 任务
  participant D as drive 与运行时
  C->>L: TCP 连接
  L->>L: 先 try_acquire live 许可,再取 handshake 许可
  L->>S: open(tag, peer, 计数用 wire, stop)
  S-->>L: carrier 会话,token 是根的子 token
  L->>K: spawn_until(会话 token, serve_socket)
  K->>T: accept(tcp, sink)
  T->>C: 在 TRANSPORT_HANDSHAKE_TIMEOUT 内完成 TLS 握手
  T->>V: sink(stream) 在同一会话上 spawn serve_connection
  V->>D: 基于 VlessCore 的 ProxyServerRuntime
  C->>D: 请求头
  Note over D: 会话上第一个被准入的流绑定它的用户
  Note over D: 协议核心已建立,drive drop 它的 handshake 许可
  D->>D: 中继,把传输层字节累加到会话的 wire 上
  Note over V: 任务结束,drop 它的许可和会话
supervisor/src/system/listener.rs
pub(crate) async fn bind(tag: &str, bind: &BindSpec) -> io::Result<Bound>;

对每个在其绑定上还没有运行中监听器的期望入站,Actor::prepare 调用 bind,并把错误映射为 ApplyError::Bind { inbound, bind, source },显示为 inbound <tag>: binding <bind> failed: <error>。tag 只用于日志。

BindSpec bind 做什么 返回
Tcp { host, port } std::net::TcpListener::bind((host, port))、set_nonblocking(true)、tokio::net::TcpListener::from_std Bound::Stream(StreamListener::Tcp)
Udp { host, port } std::net::UdpSocket::bind((host, port)) Bound::Udp
Unix(path) 见下文 Bound::Stream(StreamListener::Unix)
Tun(Create(spec)) etemenanki_protocols::tun::open(spec).await,然后 try_clone 得到 first。在 info 级别记录 inbound <tag> owns tun device <name>。 Bound::Tun
Tun(Fd(supplied)) etemenanki_protocols::tun::adopt(supplied.fd.as_fd()),得到一个非阻塞的副本;然后 try_clone 得到 first。不创建设备,也不配置地址和路由,所以不需要特权。在 info 级别记录 inbound <tag> serves a supplied tun device。 Bound::Tun

host 不必是 IP 字面量。两个 bind 调用都通过 ToSocketAddrs 接受 (host, port),它用系统解析器解析名字,然后绑定第一个能成功绑定的解析地址;全部失败时,返回最后一个地址的错误。

tun::open(protocols/src/tun/device.rs)负责设备相关的工作:

  1. check_platform(spec):在 Linux 以外的平台上,带路由的 spec 被拒绝,报错 tun routes are installed only on Linux; add them with the OS route tool。校验对 TunSource::Create 调用同一个函数,所以 --test 与实际启动的结论一致。
  2. 用 MTU、名字(如果给了)、第一个 IPv4 地址和所有 IPv6 地址(ipv6_tuple)构建设备;其余每个 IPv4 地址都添加到已建好的设备上。
  3. 把它设为非阻塞,并读回内核最终确定的名字和接口索引。
  4. 在 Linux 上,通过 rtnetlink 把每条路由作为该接口上的 link 作用域路由安装。路由从不显式删除:它们属于该接口,最后一个描述符关闭时,内核连同接口一起删除它们。

任何失败都会 drop 构建到一半的设备,从而销毁它。副本与原描述符共享文件状态标志,所以 adopt 也会让传入的描述符变成非阻塞。在 Android 和 iOS 上,tun::open 失败并报错 tun devices are created by the system VPN API here; adopt its descriptor。设备一侧见 TUN。

  1. symlink_metadata(path) 查看该路径上有什么,不跟随符号链接:
    • 如果是 socket,就认为它是崩溃的进程留下的陈旧文件,并删除它。同一 supervisor 在同一路径上的存活监听器具有相同的绑定,会被保留,所以 bind 永远不会遇到自己的 socket;
    • 其他任何东西(包括符号链接)都是别人的文件:bind 失败,返回 io::ErrorKind::AlreadyExists 和 <path> exists and is not a socket;
    • NotFound 则继续;其他错误直接返回。
  2. std::os::unix::net::UnixListener::bind(path) 创建 socket 文件。
  3. 再次调用 symlink_metadata(path),把新文件的设备号和 inode 记录到 SocketFile 中。如果这一步失败,删除该文件并返回错误。
  4. set_nonblocking(true) 和 tokio::net::UnixListener::from_std。这里失败会 drop SocketFile,从而删除文件。

etemenanki-app 能展示两种结果。把 listen 设为一个存放着普通文件的路径时,--test 会通过(因为试运行不绑定任何东西),而启动会失败:

ERROR etemenanki_app: failed to start: inbound in: binding unix:/run/etemenanki/socks.sock failed: /run/etemenanki/socks.sock exists and is not a socket

路径空闲时,启动会记录监听器日志,收到 SIGTERM 时再删除该文件:

INFO etemenanki_supervisor::system::listener: inbound in listening on unix:/run/etemenanki/socks.sock

Unix 监听器没有本地 IP,也没有对端 IP。它不带传输层:校验拒绝其上除纯 TCP 以外的任何形态(a unix socket carries no transport; its shape must be plain tcp);在其上开启 udp 的 SOCKS 需要设置 udp_bind(socks over a unix socket has no local IP for UDP associate; set udp_bind or turn udp off)。经由 Unix socket 的 UDP ASSOCIATE 会接收哪些来源的数据报,见 SOCKS。

supervisor::check(spec) 在一个全新的 actor 上以 bind = false 运行 Actor::prepare:每个 handler、用户表、出站和路由都会构建,但不绑定、也不配对任何新监听器,所以不占用端口、不打开 Hysteria 2 endpoint,也不创建 TUN 设备。除了无法绑定的端口或设备之外,启动时会拒绝的,它都会拒绝。etemenanki-app 的 --test 运行它,成功时打印 Configuration OK.,否则在 error 级别记录 configuration invalid: <error> 并以失败状态退出。

supervisor/src/system/listener.rs
impl Pending {
pub(crate) fn pair(bound: Bound, handler: InboundHandler) -> io::Result<Self>;
}
Bound + InboundHandler pair 构建什么
Stream + Stream 不再构建其他东西:Pending::Stream(listener, inbound)
Udp + Hysteria2 check_obfs(Salamander::new 会拒绝的 Salamander 密钥在这里失败,而不是在提交时);circuits 取 connection.circuit_permits;Hy2Inbound::new(ListenerConfig { connection, quic, obfs, max_connections });inbound.open(socket) 在该 socket 上构建 QUIC endpoint。
Tun + Tun handler.inbound.prepare(first):把副本注册到 Tokio 运行时,并配置 IP 协议栈(MTU、UDP 超时,以及 macOS 和 iOS 上的 packet information)。
其他任何组合 io::ErrorKind::InvalidInput,the handler is not of the kind its listener serves

已打开的 endpoint 在开始服务之前不接纳任何人。Hy2Inbound::open:

  1. 把入站 obfs 的 Salamander 密钥(没有就是空)存入入站共享的 SalamanderSwitch,使 endpoint 从第一个包起就进行混淆。遇到 Salamander::new 拒绝的密钥(短于 MIN_PSK_LEN,即 4 字节)时 open 失败,不过 pair 已经用 check_obfs 拒绝过这种密钥;
  2. 以 gate = RecvGate::default()(一个关闭的闸门)调用 endpoint::from_socket(socket, None, switch, gate)。from_socket 把 socket 设为非阻塞,为 quinn 的 Tokio 运行时包装它,并且即使没有密钥也再用 SalamanderSocket::switchable 包装一层,这样混淆可以在运行中的 endpoint 之下开启、关闭或换成另一个密钥。没有服务端配置时,endpoint 拒绝握手;闸门关闭时,它不从 socket 读取任何东西,所以客户端的包留在内核中等待,就像等待一个还没人服务的 socket 一样。run 用 RecvGate::open 打开闸门。

同样,PreparedDevice 在其运行时启动之前不被任何东西读取。TunInbound::prepare 通过 TunDevice::new 把描述符注册到 reactor,并设置协议栈的 IpStackConfig。其 mtu setter 拒绝小于 1280 的值,这就是校验的 MIN_TUN_MTU 为 1280 的原因;prepare 把这种拒绝转成 InvalidInput。未经服务就被 drop 时,已打开的 endpoint 和已准备的设备都会关闭,入站保持原样。

经过校验之后,这种不匹配不会发生,因为校验把每种协议与它的那一种绑定配对;这项检查保证无论调用方传入什么,start 都不会失败。这个错误和 pair 的其他所有失败(密钥被拒绝、from_socket 或 quinn 创建 endpoint 失败、TunInbound::prepare 失败)一样,都成为 ApplyError::Build。

supervisor/src/system/listener.rs
impl Listener {
pub(crate) fn start(bind: BindSpec, pending: Pending, shared: &Shared, stop: CancellationToken) -> Self;
}

Actor::commit 在 plane 发布、被移除的用户撤销之后,对每个待启动配对以 stop = root.child_token() 调用它,并把监听器追加到 Actor::listeners。

Pending start 在任务跟踪器上 spawn 的任务 ListenerKind
Stream run_stream_inbound(listener, rx, shared, stop),其中 rx 是 watch::channel(inbound) 的接收端 Stream(tx)
Hysteria2 run_hysteria_inbound(inbound.clone(), tag_cell, endpoint, shared, stop),其中 tag_cell = Arc::new(ArcSwap::from_pointee(tag)) Hysteria2 { inbound, tag: tag_cell, circuits, obfs }
Tun 通过 spawn_tun 调用 run_tun_inbound(handler, prepared, rx, plane, flows, stop),其中 rx 是一个无界 mpsc 通道的接收端 Tun { device, restart: tx }

每次启动都在 info 级别记录 inbound <tag> listening on <bind>,target 为 etemenanki_supervisor::system::listener。

绑定已在运行的期望入站,要么保持不变,要么得到一个新 handler(Step::SwapHandler);保持不变的入站仍可能得到一张新用户表。规划的规则见规划与应用变更。监听器一侧:

supervisor/src/system/listener.rs
impl Listener {
pub(crate) fn prepare_swap(&self, handler: InboundHandler) -> io::Result<Swap>;
pub(crate) fn swap(&mut self, swap: Swap, shared: &Shared);
pub(crate) fn prepare_users(&self, table: UserTable) -> io::Result<Users>;
pub(crate) fn store_users(&self, Users(table): Users);
pub(crate) fn circuits(&self) -> Option<&Arc<Semaphore>>;
}
监听器 prepare_swap(准备) swap(提交)
流式 拒绝其他类型的 handler 设置 tag,然后在 watch 发送端上调用 send_replace(inbound)。accept 循环在每次 accept 时读取 handler,所以下一个 socket 得到新的 handler。
Hysteria 2 拒绝其他类型;对新密钥调用 check_obfs 如果 obfs 不同:调用 set_obfs(它的 expect 说明 prepare_swap 已检查过密钥)。然后依次更新 tag 和 tag 单元、记录的 circuit 预算,并调用 set_connection_config、set_quic 和 set_max_connections。
TUN 拒绝其他类型;复制监听器的设备描述符,并对副本调用 prepare 设置 tag,并在 restart 上发送 handler 和已准备的副本。如果运行时任务已经结束导致发送失败,spawn_tun 用它们启动一个新的运行时任务,并替换发送端。

a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections 用同一端口上从 SOCKS 改为 HTTP 的入站固定了流式替换的行为:应用之前打开的中继继续中继,新客户端说 HTTP,而在应用之前刚被接受、在应用之后才握手的 socket 仍然按 SOCKS 应答。

Hysteria 2 替换能影响到什么,由 Hy2Inbound 的各个 setter 决定:

Setter 影响范围
set_connection_config 之后到达的连接。
set_quic 新握手:运行中的 endpoint 立即采用新的服务端配置。存活连接保留它们协商好的内容。
set_max_connections 从现在起到达的连接。调低它不会关闭任何存活连接;在足够多的连接结束之前,新连接会被拒绝。这个上限是每次有连接到达时检查的计数,而不是信号量,所以可以在存活连接之下修改。
set_obfs 从它返回的那一刻起,双向的每个包。在旧密钥下协商的每个连接都会被关闭,因为它们无法在新密钥下继续;endpoint 保持绑定并继续接受连接。设置为正在使用的同一密钥不会关闭任何连接。

当该绑定上运行中的入站和期望的入站都是 Hysteria 2 且 max_circuits 相同时(supervisor.rs 中的 same_circuits),circuit 预算会被沿用:Actor::prepare 把 Listener::circuits() 传给 build_handler,所以新的连接配置携带运行中的信号量,存活的 circuit 继续占用它的额度。大小不同时则得到一个新的信号量,即一份全新的预算,与仍在旧信号量下持有的 circuit 并存。

Hysteria 2 混淆的变化,以及 TUN 入站设备或设置的变化,属于 Step::Disrupt,因为它们会结束存活连接。这样的中断性变更需要 ApplyOptions::allow_disruptive:没有它的应用会被拒绝,返回 ApplyError::Disruptive。

etemenanki-app 拒绝这样的热重载,并记录 reload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply it;重启即可应用。见etemenanki-app:运行、重载与关停。

katana 应用 spec 时设置了 allow_disruptive,所以混淆的变化在 katana 中会被应用;katana 没有 TUN 入站。

Hysteria 2 入站的伪装响应、证书、UDP 或连接数上限的变化是普通替换,不是 Step::Disrupt(a_hysteria2_masquerade_change_swaps_without_disrupting)。

监听器 prepare_users 接受 store_users
流式 fits handler 协议的表 对 handler 的协议调用 store_into
Hysteria 2 UserTable::Hysteria2 Hy2Inbound::set_authenticator(Arc::new(authenticator)),存入当前的连接配置:从现在起认证的连接使用它
TUN UserTable::None 什么都不做

其他任何表都以 the handler is not of the kind its listener serves 被拒绝。应用路径和用户编辑路径(set_users、upsert_user、remove_user)最终都走到这里。两者都要等这次变更的每张表都准备好之后才存入其中任何一张,并且都先撤销被移除用户的 principal,所以针对一张即将被替换的表进行的握手,会发现自己的 principal 已被撤销。这个顺序及其对存活会话的影响,见用户、principal 与会话。

在一次应用中,一个入站要么得到替换,要么得到用户表,不会两者兼有:有变化的入站,其 handler 在构建时已经带上了新用户;只有未变化但准入结果不同的入站才会得到一张表。所以 set_connection_config 和 set_authenticator 永远不会在同一个 Hysteria 2 监听器上竞争;又因为 actor 一次只执行一条命令,其他调用方也不会与它们竞争。

supervisor/src/serve.rs
const ACCEPT_ERROR_BACKOFF: Duration = Duration::from_millis(100);
const MAX_HANDSHAKES_PER_INBOUND: usize = 204_800;
const MAX_LIVE_CONNECTIONS_PER_INBOUND: usize = 4_194_304;
pub(crate) async fn run_stream_inbound(
listener: StreamListener,
handler: watch::Receiver<Arc<StreamInbound>>,
shared: Shared,
stop: CancellationToken,
);

循环是一个 biased 的 tokio::select!,stop.cancelled() 优先于 listener.accept()。biased 表示每一轮都先检查停止信号,所以监听器一旦停止,就不会再多接受一个 socket。停止时循环退出并 drop 监听器,从而关闭端口或删除 socket 文件。

supervisor/src/serve.rs
fn should_backoff_accept_error(e: &io::Error) -> bool;
async fn backoff_or_cancelled(token: &CancellationToken) -> bool;
io::ErrorKind 含义 处理
ConnectionAborted、Interrupted 某个客户端在 accept 之前离开,或者信号中断了调用。监听器本身没有问题。 在 debug 级别记录 accept error: <error>,然后立即再次 accept
其他 通常是资源耗尽,例如文件描述符用完,立即重试还会再次失败 在 warn 级别记录 accept error, backing off 100ms: <error>,然后 backoff_or_cancelled 休眠 ACCEPT_ERROR_BACKOFF

退避能防止持续性错误把循环变成刷满日志的忙等。backoff_or_cancelled 让休眠与停止 token 竞争,token 先触发时返回 true,循环随即退出,不必睡完。accept 错误本身永远不会结束循环。

两个信号量在 run_stream_inbound 启动时创建,所以它们的计数属于监听器,与监听器同生命周期。替换 handler 会保留它们,因为循环一直在运行。

信号量 容量 获取时机 释放时机
live MAX_LIVE_CONNECTIONS_PER_INBOUND(4,194,304) accept 时第一个获取 连接结束时
handshakes MAX_HANDSHAKES_PER_INBOUND(204,800) accept 时第二个获取 对由 drive 服务的协议,协议核心报告已建立后,drive drop 它持有的许可句柄

两者都用 try_acquire_owned 获取:循环从不等待许可。任一个耗尽时,循环在 debug 级别记录日志并继续,这会 drop 刚接受的 socket(客户端看到连接被关闭)以及为它取得的任何许可:

dropping inbound connection; live connection limit reached
dropping inbound connection; handshake limit reached

立即拒绝能让循环保持响应,也能防止突发流量堆积起没人服务的 socket。live 上限是护栏,不是配额:繁忙的入站承载着数十万个大多空闲的连接,所以这个上限远高于正常流量。它也高于典型的单进程文件描述符上限,所以不要指望它在描述符耗尽之前起作用。描述符耗尽表现为 accept 错误,由上文的 100 ms 退避处理。面向运维的限制见限制与超时。

取得两个许可后,循环:

  1. 读取当前 handler:handler.borrow().clone(),它只在克隆 Arc 期间持有 watch 通道的读锁;
  2. 打开该 socket 的会话:shared.sessions.open(inbound.tag, peer, Wire::counted(), &stop)。从这一刻起,会话已登记,并会被 Supervisor::sessions 列出;在它的第一个流绑定用户之前,它没有用户。因此即使还在传输层握手阶段,它也会随入站一起被关闭。open 在注册表的锁内检查停止 token,token 已取消时返回 None,循环随即退出并 drop 该 socket。正是这项检查保证了对被移除入站的会话清扫不会漏掉任何会话(a_stopped_listener_opens_no_session);
  3. 用 handler、对端、作为 carrier 的会话、stop 和 Shared 的克隆以及两个许可构建一个 Connection;
  4. 用 shared.spawn_until(carrier.token(), …) spawn connection.serve_socket(socket)。

对端 IP 在 accept 时被保留而不是丢弃,因为路由需要它:没有它,RouteMatch::SourceCidr 就无从匹配;它也是 SOCKS UDP ASSOCIATE 唯一接收其数据报的客户端。入站 tag 以同样的方式供给 RouteMatch::InboundTag。两者都通过 FlowContext 到达路由(plane:为每个流选路)。

serve_socket 在被接受的 socket 上运行传输层,并服务传输层产出的每条字节流。

  • AcceptedSocket::Unix:serve_connection 直接在 Unix 流上内联运行,local_ip = None,connector 属于 carrier 会话,source = None。不运行传输层,也不设置 keepalive。
  • AcceptedSocket::Tcp:local_ip 取自 tcp.local_addr()。SOCKS 驱动把它放进应答中,并且除非设置了 udp_bind,否则把 UDP ASSOCIATE 中继绑定在它上面。然后运行 InboundTransport::accept(tcp, sink)。
protocols/src/transports/accept.rs
pub const TRANSPORT_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
impl InboundTransport {
pub async fn accept<F>(&self, tcp: TcpStream, mut sink: F) -> io::Result<()>
where
F: FnMut(Accepted);
}

accept 在 socket 上启用 TCP keepalive,在 TRANSPORT_HANDSHAKE_TIMEOUT 内完成传输层自身的握手,并为每条字节流调用 sink:

InboundTransport 字节流
Tcp 一条,即 socket 本身
Tls 一条,在 TLS 握手之后
Ws 一条,在 WebSocket 升级之后(配置了 TLS 时在 TLS 之上)
Grpc 隧道服务路径上的每条 HTTP/2 stream 各一条

超时以 io::ErrorKind::TimedOut 失败,文本为 tls handshake timed out、websocket handshake timed out 或 grpc handshake timed out。任何传输层错误都会让 serve_socket 结束,并在 debug 级别记录 inbound transport failed: <error>,socket 和它的会话随之被 drop。传输层见传输层:TCP 与 TLS和传输层:WebSocket 与 gRPC。

sink 闭包决定每条字节流属于哪个会话:

  • 每个 socket 一条流(TCP、TLS、WebSocket)。 这条流作为 carrier 会话来服务:connector_for(tag, peer, carrier),serve_connection 在 carrier 的 token 之下 spawn。
  • gRPC。 每条 HTTP/2 stream 都是一个独立的会话,通过 Shared::connector 打开,有自己的 Wire。

serve_connection:每种协议一个驱动

Section titled “serve_connection:每种协议一个驱动”
supervisor/src/serve.rs
async fn serve_connection<S>(
inbound: Arc<StreamInbound>,
stream: S,
local_ip: Option<IpAddr>,
connector: AppConnector,
live: Arc<OwnedSemaphorePermit>,
handshake: Arc<OwnedSemaphorePermit>,
) where
S: AsyncRead + AsyncWrite + Unpin + Send + 'static;

serve_connection 从 connector 读取 source(对端 IP;经由 Unix socket 时为 None),从 handler 读取 sniff,然后按协议分派:

StreamProtocol 驱动 协议核心 运行时 BUF
Socks 在 socks.load_full() 上调用 SocksInbound::serve,外包 WireCounted 无 不适用
Http drive HttpCore::new(config.load_full(), sniff, source) HttpCore::<Principal>::BUF_SIZE = MAX_HEAD = 65,536
Trojan drive TrojanCore::new(validator.load_full(), sniff, source) 16,384
Vless drive VlessCore::new(validator.load_full(), sniff, source) 16,384
Vmess drive VMessCore::new(validator.clone(), now_unix, sniff, source) 32,768
Shadowsocks drive ShadowsocksCore::new(resolved.load_full(), sniff, source) 20,480
Ss2022 drive 先 users.load_full(),再 Ss2022Core::with_system_clock(users.config, users.validator, sniff, source) MAX_RECORD_LEN = 65,569

每个协议核心都以 Principal 实例化,每个 BUF_SIZE 都是协议核心的关联常量,作为运行时的 const 泛型参数传入,让每种协议得到按其最大帧确定大小的缓冲区。Shadowsocks 2022 的值是对端可能发送的最大记录:MAX_PACKET_SIZE(65,535)加上 RECORD_OVERHEAD(2 字节长度和两个 16 字节 tag,共 34 字节)。运行时为每个连接分配三个 BUF 字节的缓冲区(传输层读取、传输层暂存、出站临时区),所以这些常量决定了一个连接的大部分固定内存:HTTP 为 196,608 字节,Shadowsocks 2022 为 196,707 字节,VMess 为 98,304 字节,Shadowsocks 为 61,440 字节,Trojan 和 VLESS 为 49,152 字节。缓冲区见服务端运行时。

无论驱动返回什么,错误都会连同协议名和客户端地址在 debug 级别记录,Ok 则不记录:

vless connection from Some(203.0.113.7) ended: inbound handshake timed out after 10s
socks connection from None ended: client did not complete its request in time

SOCKS 套不进协议核心的模型:它的 UDP 一侧在第二个 socket 上,控制连接只负责维持它存活。因此它有自己的驱动:

protocols/src/socks/server.rs
impl<T: Send + Sync + 'static> SocksInbound<T> {
pub async fn serve<S, C>(
&self,
stream: S,
local_ip: Option<IpAddr>,
source: Option<IpAddr>,
connector: C,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
C: Connector<Flow<T>>,
C::Datagram: DatagramLink<Addr = Destination>;
}
  • 驱动用它自己的 tokio::time::timeout(HANDSHAKE_TIMEOUT, …) 限制整个握手,超时报错 client did not complete its request in time。之后它中继 CONNECT,或驱动 UDP ASSOCIATE;见 SOCKS。
  • 没有运行时来报告这条流传输了多少字节,所以 serve_connection 用会话的 wire 把它包进 WireCounted,SOCKS 的问候、认证和请求像其他字节一样计入会话。
  • connector 以 connector.charging_datagrams() 的形式传入。UDP 关联的数据报从不经过控制流,所以改为把其子 link 收发的每个数据报的载荷加到会话的 wire 上(按用户的用量计费)。
supervisor/src/serve.rs
pub trait Established {
fn is_established(&self) -> bool;
}
type Step = Result<Traffic, RuntimeError<io::Error>>;
async fn drive<const BUF: usize, Core, S>(
stream: S,
core: Core,
connector: AppConnector,
handshake: Arc<OwnedSemaphorePermit>,
) -> io::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
Core: ProxyCoreDecode<Target = Flow, Error = io::Error, TransportAddr = ()> + Established;

drive 构建 ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector).showing_progress(),这是一个每完成一个工作单元就产出一个 Traffic 增量的 Stream。drive 需要这些产出:在两次产出之间,它可以询问 runtime.core().is_established(),并在恰当的时刻停止监视。

stateDiagram-v2
  [*] --> Handshake
  Handshake --> Handshake: Some(Ok),已计数,协议核心未建立
  Handshake --> Relay: 协议核心已建立,drop handshake 许可
  Handshake --> TimedOut: HANDSHAKE_TIMEOUT 内没有步骤
  Handshake --> Failed: Some(Err)
  Handshake --> Ended: None
  Relay --> Relay: Some(Ok),已计数
  Relay --> Failed: Some(Err)
  Relay --> Ended: None
  TimedOut --> [*]
  Failed --> [*]
  Ended --> [*]
  1. 握手(Handshake)。 协议核心尚未建立时,每次 runtime.next() 都包在 tokio::time::timeout(HANDSHAKE_TIMEOUT, …) 中。
    • Some(Ok(traffic)) 被计数,循环继续。
    • 超时返回 io::ErrorKind::TimedOut,文本为 inbound handshake timed out after 10s(由 HANDSHAKE_TIMEOUT.as_secs() 拼出)。
    • Some(Err(e)) 返回 io::Error::other(e.to_string()),所以日志显示运行时的文本。RuntimeError 给传输层失败加上前缀 transport: ,给协议核心错误加上前缀 proxy core: (例如 proxy core: client did not complete its request in time);它的其他变体是运行时自己的文本:effect targets an unknown outbound key、open reuses a live outbound key、forward range outside the event slice or held buffer、core consumed an impossible byte count、protocol frame exceeds the read buffer、stream effect on a datagram outbound or vice versa 和 bytes staged toward a datagram transport without a peer。它们各自何时出现,见服务端运行时。
    • None 表示运行时在建立之前就结束了,例如客户端关闭了连接,或者协议核心拒绝了请求并结束。drive 返回 Ok(())。
  2. 中继(Relay)。 协议核心一旦建立,drive 就 drop 它的 handshake 许可(对存放许可的 Option 调用 handshake.take()),然后一直 poll 运行时直到结束,同样对每一步计数,并以同样的方式转换 RuntimeError。

计数供给会话:在两个阶段中,每一步的 transport_rx 都作为上行、transport_tx 作为下行加到会话的 Wire 上。这些是传输层之后协议流的字节,包括握手。它们对计费意味着什么,见按用户的用量计费。

每个流式协议核心都根据自己的 Timing(protocols/src/core/mod.rs)回答 is_established:在 Phase::Relay 和 Phase::Closing 中为真。协议核心在请求解析完毕、它的流正在打开时进入 Phase::Relay,所以“已建立”表示请求已被接受,而不是出站已连接。仍在收集嗅探前缀的协议核心(Phase::Sniff,受 SNIFF_TIMEOUT 即 300 ms 限制)尚未建立,所以看门狗也覆盖嗅探阶段。

客户端发送数据之前,运行时不会向协议核心投递任何事件,而协议核心在收到第一个字节事件时才设置它的握手截止时间。因此,一个连上来却一言不发的客户端会永远停留在运行时中。drive 从外部堵住这个缺口,两个计时器互为补充:

计时器 何时设置 度量什么 捕获什么
drive 的 tokio::time::timeout 建立之前的每次 runtime.next() 两个运行时步骤之间的时间 从不说话或中途沉默的客户端
协议核心的 Timing 截止时间 第一个字节事件,只设一次 自请求开始以来的时间 不停发送却始终不完成请求的客户端

两者都使用 protocols/src/core/mod.rs 中的 HANDSHAKE_TIMEOUT,即 10 秒。

is_established 是每个协议核心上的固有方法(inherent method),而泛型函数无法调用固有方法。serve.rs 声明了 Established,并用 established! 宏为 HttpCore<Principal>、TrojanCore<Principal>、VlessCore<Principal>、VMessCore<Principal>、ShadowsocksCore<Principal> 和 Ss2022Core<Principal> 实现它,转发到各自的固有方法。要让 drive 服务一个新的协议核心,必须把它加到这个列表里。

只要绑定还在 spec 中,Hysteria 2 入站就一直持有它的 UDP socket 和 QUIC endpoint。QUIC 连接内部的一切,从凭据交换到代理流和数据报,都属于 protocols crate;见 Hysteria 2:服务端。本节介绍 supervisor 在它周围做的事。

supervisor/src/serve.rs
pub(crate) async fn run_hysteria_inbound(
inbound: Hy2Inbound<Principal>,
tag: Arc<ArcSwap<CompactString>>,
endpoint: OpenedEndpoint,
shared: Shared,
stop: CancellationToken,
);
protocols/src/hysteria/server/inbound.rs
impl<T> Hy2Inbound<T> {
pub fn open(&self, socket: std::net::UdpSocket) -> io::Result<OpenedEndpoint>;
pub async fn shutdown(&self);
}
impl<T: Send + Sync + 'static> Hy2Inbound<T> {
pub async fn run<C, F>(&self, endpoint: OpenedEndpoint, make_connector: F, token: CancellationToken)
where
F: Fn(IpAddr, ConnectionBytes) -> (C, CancellationToken) + Send + Sync + 'static,
C: Connector<Flow<T>> + Clone + Send + 'static,
C::Future: Send,
C::Stream: Send,
C::Datagram: DatagramLink<Addr = Destination> + Send;
}
#[derive(Clone)]
pub struct ConnectionBytes(quinn::Connection);
impl ConnectionBytes {
pub fn bytes(&self) -> (u64, u64);
}

在准备阶段,build_handler 用证书和私钥构建 QUIC 服务端配置,listener::bind 绑定 UDP socket,Pending::pair 在其上打开 endpoint,所以其中任何一步失败,都会在任何东西发生变化之前拒绝这次应用。提交时启动的任务开始运行后,run 安装 QUIC 服务端配置,并在入站的锁内记录该 endpoint。所以,在此之前调用的 set_quic 提供的就是 run 安装的配置,在此之后调用的则应用到运行中的 endpoint。然后 run 打开接收闸门,在一个 biased 的 tokio::select! 上循环,优先级依次为:

  1. token(只响应一次):取消它会让 run 转入排空;
  2. 一个已结束的连接任务,这样在新客户端涌入之前先回收已结束的任务;
  3. 下一个到来的客户端。endpoint.accept() 产出 None 表示 endpoint 已关闭,循环结束。

每个客户端通过一个 Admitted 守卫在 max_connections 之下占一个位置:对一个共享计数做 fetch_update,只有在已占位置少于上限(在该客户端到达时读取)时才成功。守卫移入连接的任务,任务结束时归还位置。找不到位置,或在 run 排空期间到达的客户端,会收到 QUIC 拒绝。被准入连接的任务还会记录它到达时当前的混淆 token,set_obfs 正是靠它精确地关闭那些在被替换密钥下协商的连接。正在排空的 run 在最后一个连接任务结束后返回;endpoint 先关闭的 run 则等剩下的任务结束后返回。

每个连接的 QUIC 握手一完成,run 就用客户端 IP 和一个 ConnectionBytes 调用一次 make_connector。supervisor 的闭包:

  1. 从 tag 单元读取入站当前的 tag。跨应用保留的监听器可能在服务一个改了名的入站,所以 tag 按连接读取;
  2. 用 Shared::connector(tag, Some(ip), Wire::quic(bytes), &stop) 打开一个会话,返回 connector 和会话的 token。该连接的每条代理流和它的数据报通道各拿到一份 connector 的克隆,所以它们同属一个会话;
  3. 如果客户端还在握手时监听器已被停止,就返回一个没有会话的 connector 和一个已取消的 token。连接会立即被关闭,所以没有会话能逃过被移除入站的关闭。

取消返回的 token 只会关闭这一个连接,使用 QUIC 应用错误码 0x107(DISCONNECT_CODE,即 HTTP/3 的 H3_EXCESSIVE_LOAD)和原因 disconnected by the server,不影响其他任何东西。因混淆变化而关闭的连接得到同样的错误码,原因为 obfuscation changed。

ConnectionBytes 持有 quinn::Connection,从它的统计中读取 (udp_rx.bytes, udp_tx.bytes):该连接的每个 UDP 数据报,包括握手、QUIC 分帧和 TLS。每当读取会话的字节数时,Wire::quic 就读取它,所以 Hysteria 2 会话统计的是它的 QUIC 连接的字节,而不是它承载的载荷(a_hysteria2_session_bills_its_quic_connections_bytes)。因为它持有该连接,只要会话的 wire 还持有它,它就能一直读到该连接的计数。

run_hysteria_inbound 运行 inbound.run(endpoint, make, stop),并让它与会话注册表的根 token 竞争:

  • run 先返回。 取消 stop 让 run 进入排空:之后到达的客户端被拒绝,已在服务的连接继续运行,直到它们结束或通过会话被关闭。最后一个结束后,run 返回,inbound.shutdown() 释放端口。
  • 根 token 先被取消。 这只发生在 Supervisor::shutdown 的宽限期结束之后。仍在 QUIC 握手中的客户端还没有会话,根 token 的取消触及不到它,所以任务同时运行 run 和 inbound.shutdown():shutdown 在正在排空的 run 之下关闭 endpoint,立即结束其上的每个连接和握手,run 也随之返回。shutdown_does_not_wait_out_a_stalled_hysteria2_handshake 固定了“卡住的握手不会拖住关停”这一行为。

Hy2Inbound::shutdown 以应用错误码 0x100(CLOSE_CODE)关闭 endpoint,并在有限时间内等待端口释放;没有运行中的 endpoint 时立即返回。它的步骤和超时见 Hysteria 2:服务端。

TUN 入站持有一个三层设备。protocols crate 中的用户态 IP 协议栈终结操作系统引入该设备的每条 TCP 连接和 UDP 流;见 TUN。

supervisor/src/serve.rs
pub(crate) async fn run_tun_inbound(
handler: TunHandler,
prepared: PreparedDevice,
restart: mpsc::UnboundedReceiver<(TunHandler, PreparedDevice)>,
plane: PlaneCell,
flows: Tracker,
stop: CancellationToken,
);

监听器保留设备自己的描述符,每个运行时拿到的是一个副本,这个副本在启动该运行时的那次变更提交之前就已准备好。run_tun_inbound 循环执行:

  1. 构建 connector 工厂:对每个流调用 AppConnector::new(plane, FlowContext { inbound_tag: tag, source: Some(ip) }, flows, None),其中 ip 是该流的源地址。
  2. 在 token = stop.child_token() 之下运行 handler.inbound.run(prepared, make, token),并让它与 restart.recv() 竞争。
  3. 如果收到重启请求,取消 token 并等待 run 返回。如果 run 自行返回(因为 stop 已触发或设备已消失),就没有下一个运行时。
  4. handler.inbound.shutdown().await:每隔 RELEASE_POLL(20 ms)轮询一次,最多等待 RELEASE_TIMEOUT(3 秒),直到旧运行时的设备描述符已关闭且没有剩余 TCP 流。等待超时时,在 warn 级别记录 tun: device fd or <n> tcp flows still open after 3s 并继续。
  5. 如果有重启请求且 stop 未取消,就取出新的 handler 和已准备的副本,继续循环;否则返回。

所以,在下一个运行时启动之前,前一个运行时会被取消,并有最多 RELEASE_TIMEOUT 的时间释放它的描述符和 TCP 流。重启通道是无界的,但只有 actor 会写入它,每次替换该监听器的应用写入一次。监听器被 drop 时,它的发送端随之消失,recv 返回 None,循环以同样的方式结束。

替换 TUN handler 会结束设备上的流,所以规划把这种替换标记为中断性变更,与设备本身的变化一样。etemenanki-app 在热重载时拒绝 TUN 变更;重启即可应用。

TunInbound::run 服务一个已准备的设备,直到它的 token 被取消或设备消失。每个运行时都有自己的 max_flows 信号量,由 run 创建,所以重启后计数从零开始。

IP 协议栈产出 run 的处理
一条 TCP 连接 用 try_acquire_owned 取一个流许可;取不到则 drop 该连接,并在 debug 级别记录 tun: dropping a flow; the flow limit is reached。通过 serve_stream 在 PassthroughCore 之上服务它(除非 handler 关闭了嗅探,否则进行嗅探);错误在 debug 级别记录为 tun: tcp flow <src> -> <dst> ended: <error>。
udp 关闭时的 UDP 流 丢弃。
来自已有关联的源的 UDP 流 通过该关联的有界通道(FLOW_QUEUE,16)交给它;发现通道已满的流被丢弃,由客户端重传。
来自新源的 UDP 流 取一个流许可(取不到则以同样的 debug 日志丢弃),并为该源启动它唯一的关联:一个基于 TunUdpCore 的运行时。关联结束时在 trace 级别记录 tun: udp association of <src> ended: <error>。
其他任何协议,包括 ICMP 丢弃,在 trace 级别记录 tun: dropping a packet of an unsupported protocol。

前端程序填入的默认值来自 protocols/src/tun/config.rs:DEFAULT_MTU 为 1500,DEFAULT_UDP_IDLE_TIMEOUT 为 60 秒,DEFAULT_MAX_FLOWS 为 65,536。协议栈、关联和 TCP 跟踪见 TUN。

TUN 设备不准入任何用户:它的用户表是 UserTable::None,每个流都携带 Principal::anonymous(),connector 没有会话。所以 TUN 流不会被 Supervisor::sessions 列出,Supervisor::close 触及不到它们,它们也不计入任何用户的用量。这些流仍然在流跟踪器中登记和计量,只是没有会话。经过设备的 TCP 流在 protocols crate 中有自己的沉默客户端看门狗(serve_stream 在其 PassthroughCore 建立之前对每一步施加 HANDSHAKE_TIMEOUT,超时报错 tun: the client never spoke)。

Listener::stop 取消监听器的 token。

监听器 停止的部分 继续运行的部分
流式 accept 循环,它会 drop 监听器:端口关闭,Unix socket 文件被删除 每个已接受的连接,在其会话 token 之下
Hysteria 2 准入:新客户端被拒绝 每个有会话的 QUIC 连接,直到它结束或其会话被关闭。最后一个结束后,endpoint 关闭,端口释放。
TUN 运行时(它的 token 是监听器 token 的子 token):设备上的每个流随之结束 无

入站被移除或移到另一个绑定时,Actor::commit 为 Step::StopAccepting 停止监听器,并用 swap_remove 把它从 Actor::listeners 中移除。监听器随即被 drop,这同时 drop 了 TUN 设备自己的描述符、actor 持有的 Hysteria 2 句柄和流式 handler 的 watch 发送端。

连接是否也一并结束,由规划决定:

变化 步骤 连接
入站被移除 先 StopAccepting,再 CloseSessions 关闭:Sessions::close(Scope::Inbound(tag)) 取消该 tag 的每个会话(removing_an_inbound_closes_its_sessions)
入站移到另一个绑定 Bind、StopAccepting 规划不关闭它们:会话属于 tag,而 tag 没变
入站在同一绑定上改名 SwapHandler、旧 tag 的 CloseSessions 旧 tag 的会话被关闭;监听器保留

StopAccepting 在 CloseSessions 之前,而 Sessions::open 在清扫所持的同一把锁内检查停止 token,所以一个在其入站被移除时正被接受的连接,要么被清扫掉,要么根本打不开会话。

Actor::commit 在它的 ApplyReport 中报告监听器相关的步骤:

字段 来源
swapped 每个 SwapHandler
restarted 每个 Disrupt
rebound 其 tag 在运行中的 spec 里已存在于另一个绑定上的 Bind;新 tag 的绑定不列入
removed 每个 CloseSessions

a_hysteria2_obfuscation_change_needs_allow_disruptive 读取这些字段:被允许的混淆变化把该入站列在 restarted 中,而不在 rebound 中。

Supervisor::shutdown(grace) 请求 actor 执行关停,actor 会:

  1. 停止每个监听器,并取消每个负载均衡器的探测 token;
  2. 关闭任务跟踪器,最多等待 grace 让其上的每个任务结束;
  3. 取消根 token。每个会话 token 都是它的子 token,所以仍在运行的每个连接任务都会被 drop,每个 Hysteria 2 监听器立即关闭它的 endpoint;
  4. 再次等待任务跟踪器(不设时限),并清空 Actor::listeners;
  5. 停止采样器和用量 sink;此时每个会话都已结束并计入了最后的字节,它们再报告一次。

etemenanki-app 收到 SIGTERM 或 Ctrl-C 时传入 SHUTDOWN_GRACE,即 5 秒。最后一个 Supervisor 句柄在没有关停的情况下被 drop 时,actor 自己运行 shutdown(Duration::ZERO)。TUN 流没有宽限期:第 1 步就停止了它们的运行时。

移除和排空的宽限计时器运行在同一个任务跟踪器上,在根 token 取消时结束。所以一个待触发的计时器会让第 2 步等满整个 grace,而第 3 步会结束它,而不是等到它的截止时间(shutdown_ends_pending_grace_timers)。

服务端是一个 ProxyCoreDecode 协议核心、经由 TCP 或 Unix socket 提供服务的协议,需要接入以下位置:

  1. Spec。 在 supervisor/src/topology/spec_plan/inbound.rs 的 InboundProtocolSpec 中添加一个变体。如果它运行在传输层之上,把它加到 supervisor/src/build/inbound.rs 的 stream_transport 中。

  2. 校验。 supervisor/src/build/validate.rs 中的 validate_inbound 已经把 Hysteria 2 和 TUN 以外的每种协议与 TCP 或 Unix 绑定配对。添加该协议自己的检查:传输层检查;如果它没有开放模式,还要像 Trojan、VLESS 和 VMess 那样要求用户集。见校验与应用错误。

  3. 凭据。 在 credential_kind(supervisor/src/build/users.rs)中返回用来准入用户的凭据种类。

  4. 用户表。 添加一个 UserTable 变体,在 user_table 中构建它,并把这一对加到 UserTable::fits 和 UserTable::store_into 中。

  5. Handler。 在 supervisor/src/topology/inbound/mod.rs 中添加一个 StreamProtocol 变体,把表放在 ArcSwap 中(或者像 VMess 那样放在一个自己替换用户的类型中),在 StreamProtocol::name 中给它命名,并在 build_handler 中把表映射到它。

  6. 服务。 把协议核心加到 supervisor/src/serve.rs 的 established! 列表中,并添加一个 serve_connection 分支,调用 drive::<{ NewCore::<Principal>::BUF_SIZE }, _, _>,传入由表的 load_full()、sniff 和 source 构建的协议核心,再传入 connector 和 handshake 许可。drive 要求 Target = Flow、Error = io::Error 和 TransportAddr = (),并且协议核心的 is_established 只能在请求解析完毕并被接受之后才变为真;现有的协议核心用 Timing::enter(Phase::Relay, …) 做到这一点。

  7. 中断性变更。 如果某项设置的变化无法保全存活连接,把它加到 supervisor/src/topology/spec_plan/plan.rs 的 swap_disruption 中(规划与应用变更)。

  8. 前端程序。 在每个应当提供该协议的前端程序中实现它的 lowering(降为 spec):etemenanki-app 在 app/src/lower.rs 中(etemenanki-app:从 TOML 到 spec),katana 在 src/lower/inbound.rs 中(将入站与出站降为 spec)。

无法做成协议核心的协议(例如带第二个 UDP socket 的 SOCKS),需要在 serve_connection 中有自己的驱动和自己的握手超时,并使用 WireCounted 流,使它的会话仍能统计字节。

不变量 保证机制 固定它的测试
被拒绝的应用不留下任何监听器,也不服务它绑定过的任何东西。 绑定和配对发生在准备阶段;drop 一个 Pending 会关闭它的 socket、删除它的 Unix 文件、关闭它的 endpoint 或设备。只有提交阶段的 Listener::start 才提供服务。 supervisor/tests/hot_swap.rs 中的 a_refused_spec_changes_nothing
不变的绑定在应用前后保留它的 socket。 按 BindSpec 查找监听器;有变化的入站走 prepare_swap 和 swap,从不重新绑定。 a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections、a_hysteria2_obfuscation_change_needs_allow_disruptive(rebound 为空);app/tests/integration/e2e_hysteria_inbound.rs 中的 a_reload_keeps_the_inbound_serving_on_its_udp_port
保留入站及其用户的应用不触碰它的连接。 Step::Reuse:不替换;只有准入结果变化时,才给被复用的入站存入用户表(same_admissions)。 an_established_connection_survives_an_apply_that_keeps_its_inbound
新 socket 由它被接受时的当前 handler 服务。 每次 accept 调用 handler.borrow().clone();swap 只替换 watch 中的值。 a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections
停止流式监听器只停止接受,别的什么都不停。 会话 token 是根的子 token,不是监听器的子 token。 规划一侧:supervisor/tests/unit/plan.rs 中的 a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions。没有测试在有存活连接时移动监听器。
已停止的监听器上不会打开会话。 Sessions::open 在注册表锁内检查 stop;accept 循环的 select! 偏向 stop。 supervisor/tests/unit/session.rs 中的 a_stopped_listener_opens_no_session
移除入站会关闭它的会话。 StopAccepting 之后的 Step::CloseSessions。 removing_an_inbound_closes_its_sessions
Unix 监听器只删除它自己创建的 socket 文件。 SocketFile 删除前比较设备号和 inode。 删除:app/tests/integration/e2e_unix.rs 中的 socks_over_a_unix_socket_relays_and_cleans_up。设备号和 inode 检查没有专门的测试。
accept 循环从不等待许可。 try_acquire_owned;被拒绝的 socket 被 drop。 没有专门的测试
持续性 accept 错误不会让循环空转,单个客户端的错误不会拖慢循环。 should_backoff_accept_error 和 ACCEPT_ERROR_BACKOFF supervisor/tests/unit/serve.rs 中的 accept_error_backoff_classification
始终不完成请求的客户端会被断开。 建立之前 drive 按步施加的 HANDSHAKE_TIMEOUT,以及协议核心的 Timing 截止时间 协议核心一侧:protocols/tests/unit/core/mod.rs 中的 timing_arms_handshake_once_then_idle_per_byte_event,以及各协议核心的测试,例如 protocols/tests/unit/trojan/core.rs 中的 handshake_deadline_fails_the_connection。drive 自己的超时没有测试。
会话的 wire 统计其客户端一侧承载的字节。 drive 累加每一步的传输层字节;SOCKS 用 WireCounted 和计费数据报;Hysteria 2 用 Wire::quic supervisor/tests/tracking.rs 中的 mux_sub_flows_and_their_session_count_exactly_what_moved、a_udp_association_counts_each_sub_link_and_charges_its_session、a_hysteria2_session_bills_its_quic_connections_bytes
路由能看到 accept 时捕获的入站 tag 和客户端地址。 由 connector_for 和 Shared::connector 构建的 FlowContext app/tests/integration/e2e_route_context.rs 中的 inbound_tag_selects_the_route、source_cidr_matches_the_client_address
Hysteria 2 监听器永远不会在一个端口上服务两个 endpoint。 变化由各 setter 在运行中的 endpoint 内部应用。 a_reload_keeps_the_inbound_serving_on_its_udp_port;protocols/tests/pipeline/hysteria.rs 中的 switching_the_obfuscation_closes_the_old_clients_and_admits_new_ones
关停时间受宽限期约束,Hysteria 2 握手也不例外。 根 token 取消时,在正在排空的 run 之下运行 Hy2Inbound::shutdown。 shutdown_does_not_wait_out_a_stalled_hysteria2_handshake
在下一个 TUN 运行时启动之前,前一个会被取消,并有最多 3 秒释放其描述符和 TCP 流。 run_tun_inbound 在进入下一轮之前先 await run,再 await TunInbound::shutdown。 没有专门的测试
情况 位置 结果 日志
TCP 或 UDP 地址已被占用或不允许绑定 listener::bind ApplyError::Bind;这次应用不改变任何东西 前端程序的报告,例如 etemenanki-app 的 failed to start: inbound <tag>: binding <bind> failed: <OS error>
Unix 路径上是 socket 以外的东西 listener::bind ApplyError::Bind … failed: <path> exists and is not a socket
TUN 设备无法创建、配置地址或配置路由 tun::open ApplyError::Bind 操作系统错误;路由出错时为 route <net>/<prefix>: <error>
证书或私钥无法解析 build_handler ApplyError::Build building inbound <tag> failed: <error>,例如 hysteria2: the certificate file contains no certificates
handler 或用户表的类型不对 Pending::pair、prepare_swap、prepare_users ApplyError::Build(经过校验后不可达) building inbound <tag> failed: the handler is not of the kind its listener serves
短于 MIN_PSK_LEN(4 字节)的 Salamander 密钥 check_obfs ApplyError::Build;校验会先以 salamander obfs key must be at least 4 bytes 拒绝它 malformed: hysteria2 obfs psk is too short
无法在 socket 上创建 Hysteria 2 endpoint Hy2Inbound::open(from_socket、quinn 的 endpoint) ApplyError::Build building inbound <tag> failed: <error>
TUN 描述符无法注册到 reactor,或协议栈拒绝其 MTU 由 Pending::pair 或 prepare_swap 调用的 TunInbound::prepare ApplyError::Build;MTU 的情况经过校验后不可达 building inbound <tag> failed: <error>
替换时无法复制监听器的 TUN 描述符 prepare_swap(try_clone) ApplyError::Build building inbound <tag> failed: <error>
accept 以 ConnectionAborted 或 Interrupted 失败 run_stream_inbound 立即再次 accept debug accept error: <error>
accept 以其他错误失败 run_stream_inbound 休眠 100 ms;期间被停止则结束 warn accept error, backing off 100ms: <error>
live 或 handshake 信号量耗尽 run_stream_inbound socket 被 drop debug dropping inbound connection; live connection limit reached 或 … handshake limit reached
监听器在 socket 被接受的同时停止 Sessions::open 返回 None 循环结束,socket 被 drop 无
传输层握手失败或超过 10 秒 serve_socket socket 和会话被 drop debug inbound transport failed: <error>
建立之前 10 秒内没有运行时步骤 drive TimedOut debug <protocol> connection from <source> ended: inbound handshake timed out after 10s
第一个字节之后 10 秒请求仍未完成 协议核心 运行时错误 debug … ended: proxy core: client did not complete its request in time
SOCKS 握手 10 秒内未完成 SocksInbound::serve TimedOut debug socks connection from <source> ended: client did not complete its request in time
其他任何运行时或驱动错误 drive、SocksInbound::serve 连接被 drop debug <protocol> connection from <source> ended: <error>
会话的 token 被取消(见 token 与任务) spawn_until 任务的 future 被 drop 无
Hysteria 2 QUIC 握手失败 Hy2Inbound::run 连接被 drop debug hysteria2: a handshake failed: <error>
Hysteria 2 客户端在达到上限时或排空期间到达 Hy2Inbound::run QUIC 拒绝 debug hysteria2: refusing a connection; the listener is full or stopping
endpoint 关闭 3 秒后 Hysteria 2 端口仍被占用 Hy2Inbound::shutdown 照常返回 warn hysteria2: <address> did not come free within 3s
TUN 运行时的描述符或 TCP 流在 3 秒后仍然存在 TunInbound::shutdown 循环继续 warn tun: device fd or <n> tcp flows still open after 3s
TUN 替换时发现运行时任务已结束 Listener::swap 在已准备的副本上 spawn 一个新的运行时任务 无
Unix socket 文件无法删除 SocketFile::drop 忽略;下一次绑定会接管这个陈旧的 socket debug could not remove <path>: <error>

serve.rs 的日志使用 target etemenanki_supervisor::serve,监听器的日志使用 etemenanki_supervisor::system::listener。serve.rs 中关于单个连接的日志没有高于 debug 级别的。

对流式连接而言,取消就是 drop。会话的 token 触发时,spawn_until 的 select! 会在连接的 future 当前停靠的任意 .await 处 drop 它:客户端 socket 关闭,运行时及其打开的每个出站被 drop,连接的许可句柄被 drop,最后一个 Arc<Session> 句柄注销会话,并把它最后的字节计入其用户的账户。不会发送任何协议层面的告别消息。因此,协议和传输层代码必须把清理工作放在 Drop 中,绝不能放在某个 .await 之后。Hysteria 2 会话则由 protocols crate 以 DISCONNECT_CODE 关闭(见每个 QUIC 连接一个会话)。

常量 值 定义位置 作用范围
MAX_LIVE_CONNECTIONS_PER_INBOUND 4,194,304 supervisor/src/serve.rs 每个流式监听器,贯穿其整个生命周期
MAX_HANDSHAKES_PER_INBOUND 204,800 supervisor/src/serve.rs 每个流式监听器,贯穿其整个生命周期
ACCEPT_ERROR_BACKOFF 100 ms supervisor/src/serve.rs 每次需要退避的 accept 错误
HANDSHAKE_TIMEOUT 10 秒 protocols/src/core/mod.rs drive 的逐步看门狗、协议核心的 Timing、SOCKS 握手、TUN TCP 流
TRANSPORT_HANDSHAKE_TIMEOUT 10 秒 protocols/src/transports/accept.rs 入站的 TLS 握手、WebSocket 升级和 HTTP/2 握手(InboundTransport::accept)
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs 嗅探窗口,位于握手阶段之内
协议核心的 BUF_SIZE 16,384 到 65,569 字节 各协议核心 每个流式连接三个这么大的缓冲区
DEFAULT_MAX_CONNECTIONS 262,144 protocols/src/hysteria/server/config.rs 前端程序为 Hysteria 2 入站 max_connections 设的默认值
DEFAULT_MAX_CIRCUITS 4,194,304 protocols/src/hysteria/server/config.rs 前端程序为 max_circuits 设的默认值
CLOSE_CODE、DISCONNECT_CODE 0x100、0x107 protocols/src/hysteria/server/inbound.rs QUIC 应用错误码:endpoint 关闭、单个连接关闭
MIN_PSK_LEN 4 字节 protocols/src/hysteria/obfs.rs check_obfs 和 Hy2Inbound::open 接受的最短 Salamander 密钥
RELEASE_TIMEOUT、RELEASE_POLL 3 秒、20 ms protocols/src/tun/inbound.rs 每个 TUN 运行时的关停
DEFAULT_MTU、DEFAULT_UDP_IDLE_TIMEOUT 1500、60 秒 protocols/src/tun/config.rs 前端程序为 TUN 入站 mtu 和 udp_idle_timeout 设的默认值
DEFAULT_MAX_FLOWS 65,536 protocols/src/tun/config.rs 前端程序为 TUN 入站 max_flows 设的默认值,按运行时计
FLOW_QUEUE 16 protocols/src/tun/udp.rs 排队等待交给同一个 TUN 关联的新 UDP 流
MIN_TUN_MTU 1280 supervisor/src/build/validate.rs spec 可以指定的最小 TUN MTU,也是 IP 协议栈自身的下限
SHUTDOWN_GRACE 5 秒 app/src/main.rs etemenanki-app 给 Supervisor::shutdown 的宽限期

serve.rs 中的常量都不可配置。Hysteria 2 入站的 max_connections 和 max_circuits 必须至少为 1。整个工作区的限制汇总在限制、超时与内存一页上。

测试 文件 固定的行为
accept_error_backoff_classification supervisor/tests/unit/serve.rs Interrupted 和 ConnectionAborted 不退避;OutOfMemory 和 Other 退避。
a_stopped_listener_opens_no_session supervisor/tests/unit/session.rs 停止 token 已取消时,Sessions::open 返回 None,不登记任何东西。
cancelling_the_root_cancels_every_session supervisor/tests/unit/session.rs 每个会话 token 都是根的子 token。
an_established_connection_survives_an_apply_that_keeps_its_inbound supervisor/tests/hot_swap.rs 添加一个出站和一条规则的应用把该入站报告为复用,不替换也不重新绑定任何东西,存活的 SOCKS 中继继续回显。
a_protocol_change_on_the_same_port_keeps_the_socket_and_old_connections supervisor/tests/hot_swap.rs 同一端口上从 SOCKS 改为 HTTP 是替换而非重新绑定;旧中继继续工作,新客户端说 HTTP,应用之前接受的 socket 按 SOCKS 应答。
a_refused_spec_changes_nothing supervisor/tests/hot_swap.rs 一个新入站绑定失败时,同一次准备中已绑定的另一个入站被释放;plane、用户和监听器都保持原样。
removing_an_inbound_closes_its_sessions supervisor/tests/hot_swap.rs 被移除入站的连接关闭,端口拒绝连接;另一个入站的连接不受影响。
a_hysteria2_obfuscation_change_needs_allow_disruptive supervisor/tests/hot_swap.rs 默认以 ApplyError::Disruptive 拒绝;设置 allow_disruptive 后,该入站出现在 restarted 中,不在 rebound 中。
a_tun_device_change_needs_allow_disruptive supervisor/tests/hot_swap.rs 没有 allow_disruptive 时,设备变化和设置变化都以 ApplyError::Disruptive 被拒绝。需要 CAP_NET_ADMIN;无法创建设备时跳过。
shutdown_does_not_wait_out_a_stalled_hysteria2_handshake supervisor/tests/hot_swap.rs 客户端的 QUIC 握手卡住时,shutdown(Duration::ZERO) 仍在 5 秒内完成。
shutdown_ends_pending_grace_timers supervisor/tests/hot_swap.rs 长达一小时的移除和排空计时器不会拖住关停,它们暂时放过的会话被关闭。
mux_sub_flows_and_their_session_count_exactly_what_moved supervisor/tests/tracking.rs drive 向会话累加的字节数与 socket 承载的完全一致。
a_udp_association_counts_each_sub_link_and_charges_its_session supervisor/tests/tracking.rs WireCounted 统计 SOCKS 的问候、认证和请求,计费数据报累加其载荷。
a_hysteria2_session_bills_its_quic_connections_bytes supervisor/tests/tracking.rs Hysteria 2 会话每个方向统计的字节都多于 20,000 字节的载荷:统计的是 QUIC 连接的 UDP 字节。
a_protocol_change_on_the_same_bind_swaps_the_handler_and_keeps_the_listener、a_bind_change_binds_anew_and_stops_the_old_listener_without_closing_sessions、a_removed_inbound_stops_accepting_and_closes_its_sessions、an_inbound_renamed_on_the_same_bind_swaps_and_closes_the_old_tags_sessions supervisor/tests/unit/plan.rs 每种入站变化为其监听器和会话规划哪些步骤。
a_tun_settings_change_on_the_same_device_swaps_and_disrupts、a_hysteria2_obfs_change_swaps_and_disrupts、a_hysteria2_masquerade_change_swaps_without_disrupting supervisor/tests/unit/plan.rs 哪些替换属于中断性变更。
two_inbounds_may_not_share_a_bind、the_same_port_over_udp_is_another_bind、each_protocol_is_served_only_on_its_kind_of_bind、a_unix_listener_carries_only_the_plain_tcp_shape supervisor/tests/unit/validate.rs 让 BindSpec 成为 key、并把协议与绑定配对的校验。
a_connections_close_token_closes_that_connection_alone protocols/tests/pipeline/hysteria.rs 取消 make_connector 返回的 token 只关闭那一个连接,不影响相邻的连接。
cancelling_run_stops_admitting_and_waits_for_the_live_connections protocols/tests/pipeline/hysteria.rs 停止的 Hysteria 2 监听器迅速拒绝新来者,继续为存活连接中继,run 在该连接结束后返回。
a_new_connection_config_reaches_only_connections_that_arrive_after_it protocols/tests/pipeline/hysteria.rs 新的连接配置服务其后到达的连接(新密码、无 UDP)。
switching_the_obfuscation_closes_the_old_clients_and_admits_new_ones、a_too_short_obfuscation_key_is_refused_and_changes_nothing protocols/tests/pipeline/hysteria.rs 在运行中的 endpoint 上调用 set_obfs,以及它对过短密钥的拒绝。
socks_over_a_unix_socket_relays_and_cleans_up、http_connect_over_a_unix_socket_relays app/tests/integration/e2e_unix.rs 经由 Unix 监听器的一次 SOCKS 往返和一次 HTTP CONNECT;SIGTERM 之后 socket 文件消失。
a_reload_keeps_the_inbound_serving_on_its_udp_port app/tests/integration/e2e_hysteria_inbound.rs 在有客户端连接时修改 Hysteria 2 入站设置的热重载,使它继续在同一端口上服务。没有参考 hysteria 构建时跳过。
a_routed_connect_is_answered_while_the_app_runs app/tests/integration/e2e_tun.rs TUN 入站的 IP 协议栈应答一次经路由到达的 TCP 连接,app 退出后接口和路由都消失。仅限 Linux;没有 CAP_NET_ADMIN 时跳过。
inbound_tag_selects_the_route、source_cidr_matches_the_client_address app/tests/integration/e2e_route_context.rs accept 时捕获的 tag 和对端 IP 到达路由器。

在这个版本中,准入信号量和 drive 的沉默客户端超时都没有测试覆盖。修改其中任何一个时应当补上测试;暂停的 Tokio 时钟(tokio::time::pause)可以让 10 秒看门狗无需真正等待就能测试。用 cargo test -p etemenanki-supervisor 运行 supervisor 的测试,用 cargo test -p etemenanki-app 运行 app 的测试;见测试。