REST API(etemenanki-webclient)
源码文件:30 个 · 核对版本 Etemenanki 555b7df
Etemenanki/webclient/src/lib.rsEtemenanki/webclient/src/router.rsEtemenanki/webclient/src/listen.rsEtemenanki/webclient/src/etag.rsEtemenanki/webclient/src/tracking.rsEtemenanki/webclient/Cargo.tomlEtemenanki/app/src/api.rsEtemenanki/app/src/routes.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/config.rsEtemenanki/app/src/subscribe.rsEtemenanki/Cargo.tomlEtemenanki/supervisor/src/track/mod.rsEtemenanki/supervisor/src/track/sampler.rsEtemenanki/supervisor/src/entity/session.rsEtemenanki/supervisor/src/entity/id.rsEtemenanki/supervisor/src/build/apply.rsEtemenanki/supervisor/src/topology/spec_plan/mod.rsEtemenanki/supervisor/src/topology/spec_plan/plan.rsEtemenanki/supervisor/src/supervisor.rsEtemenanki/supervisor/src/topology/outbound/udp_fanout.rsEtemenanki/protocols/src/flow.rsEtemenanki/concepts/src/net.rsEtemenanki/ffi/src/proxy.rsEtemenanki/ffi/src/client.rsEtemenanki/webclient/tests/tracking.rsEtemenanki/app/tests/unit/api.rsEtemenanki/app/tests/integration/e2e_api.rsEtemenanki/ffi/tests/proxy.rs
etemenanki-webclient 是 etemenanki-app 在配置含有 [api] 段时提供的 HTTP API。它面向客户端配置,但只要配置有这个段,无论入站是什么,main 都会提供它。它列出路由组并切换它们,按请求重载,列出存活的会话和流,关闭流,并以 server-sent events 推送流事件和流量速率。路由切换的做法是编辑配置文件再重载,因此选择只存在于文件中:重启后读到的就是 API 写入的选择。
这个 crate 只负责 HTTP。路由组是什么、文件如何编辑和校验、重载做什么,都归 app 管,位于 app/src/api.rs 实现的 ConfigControl trait 之后;视图和编辑本身在 app/src/routes.rs 中,与移动端库共用。二进制程序(app/src/main.rs)绑定 API,并在重载改变 [api] 时启动、停止或重新绑定它。本页面向修改上述任一文件的贡献者。如何配置和调用 API 见用户指南。 连接相关端点读取的跟踪器见流跟踪、统计与限速,编辑所触发的重载见 etemenanki-app:运行、重载与关闭和规划并应用变更。
| 文件 | 负责 |
|---|---|
webclient/src/lib.rs |
公开类型:ConfigControl、RouteView 及其组成部分、ReloadReport 及其 JSON、ControlError、DEFAULT_GROUP;以及对其他模块的再导出 |
webclient/src/router.rs |
router():HTTP 路由、bearer 校验、edits 互斥锁,以及 ControlError 到状态码的映射 |
webclient/src/etag.rs |
ETag:配置文件字节的 SHA-256,以及 If-Match 的解析 |
webclient/src/listen.rs |
Listen、ApiOptions、OptionsError、DEFAULT_LISTEN,以及 ApiListener:它绑定 TCP 或 unix socket,并以优雅关闭的方式提供服务 |
webclient/src/tracking.rs |
连接与流量的处理函数及其 JSON 视图,直接读取 supervisor(监管器)的 Tracker |
app/src/api.rs |
options(),把 [api] 转成 ApiOptions;ConfigFile,即基于运行中 Instance 的 ConfigControl;write_atomically 和 create_beside |
app/src/routes.rs |
Snapshot:配置及其订阅文件的路由视图,以及对单个选择的编辑 |
app/src/main.rs |
start_api、Api::stop、follow_api,以及 API 在启动和关闭流程中的位置 |
webclient crate 不负责:
- 了解配置 schema。它从不解析 TOML,而是调用
ConfigControl。 - 提供 UI。它回答的是 JSON、server-sent events、
204以及 axum 自己的404和405的空 body,还有 axum 对无法解析的路径参数给出的纯文本400。 - 发送 CORS 头。router 没有添加 CORS 层。
- 提供按用户的用量。流量视图省略了采样器的按用户数据,按用户的用量只能通过 supervisor 的 Rust API 获取(按用户的用量计费)。
- 监视文件。文件监视器属于二进制程序(etemenanki-app:运行、重载与关闭);API 只在被请求时重载。
- 记录日志。它把
tracing声明为依赖,但从未调用;axum 构建时也没有启用它的tracingfeature,所以被拒绝的请求同样不会被记录。
flowchart LR app["etemenanki-app"] -->|"实现 ConfigControl,绑定 ApiListener"| web["etemenanki-webclient"] ffi["etemenanki-ffi"] -->|"RouteView、ReloadReport、ControlError"| web ffi -->|"routes::Snapshot、instance::Core"| app web -->|"Tracker 及其类型、Secret、ApplyReport"| sup["etemenanki-supervisor"] web -->|"DialNetwork、Remote"| con["etemenanki-concepts"] app --> sup
是 app 依赖 webclient crate,而不是反过来:ConfigControl 就是两者之间的接缝。这个 crate 使用了其他 workspace crate 中的以下类型:
| Crate | 类型 | 用途 |
|---|---|---|
etemenanki-supervisor |
topology::spec_plan::Secret |
bearer token |
etemenanki-supervisor |
build::apply::ApplyReport |
一次重载做了什么 |
etemenanki-supervisor |
track::{Tracker, FlowEntry, FlowEvent, StatsSnapshot, TagStats}、entity::id::FlowId、entity::session::SessionStats |
连接、流事件和流量(tracking.rs) |
etemenanki-concepts |
net::{DialNetwork, Remote} |
流的目标地址 |
移动端库复用 RouteView、ReloadReport、ControlError 和 routes::Snapshot 来实现它原生的 routes、reload 和 set_route 调用,但不提供 [api]:配置含有这个段时,它的 lowering(降为 spec)会在 info 级别记录 the config's [api] is not served here: the app reads routes and traffic through this library,并且不发布任何 API 选项。它从不调用 api::options,所以也不检查这个段:etemenanki-app 会拒绝的 [api],在那里不会被拒绝(移动端库(etemenanki-ffi))。katana 没有 REST API。
crate 本身
Section titled “crate 本身”webclient/Cargo.toml 声明了 etemenanki-webclient 版本 0.1.0,发布到一个私有 Cargo registry。它依赖 etemenanki-supervisor(版本要求 0.3.0)、etemenanki-concepts(2.0.0)、axum、tokio、tokio-stream(基于跟踪器通道的 SSE 流)、futures、compact_str、serde、sha2、hex、thiserror 和 tracing(已声明但未使用)。测试另外用到 etemenanki-protocols、tower、http-body-util 和 serde_json。
workspace 的 Cargo.toml 以 default-features = false 构建 axum 0.8.9,只启用 http1、json、query 和 tokio 这几个 feature。因此 API 只说 HTTP/1,不支持 HTTP/2。axum 的其他默认 feature,即 form、matched-path、original-uri、tower-log 和 tracing,都是关闭的;app 是唯一另一个依赖 axum 的 crate,它也没有额外启用任何 feature。
| 方法 | 路径 | 处理函数 | 成功 | 错误 |
|---|---|---|---|---|
GET |
/v1/routes |
router::list |
200,JSON 形式的 RouteView,以及文件的 ETag 头 |
422、500 |
PUT |
/v1/routes/{group} |
router::set |
200,ReloadReport,以及文件新的 ETag 头 |
400({group} 不是 UTF-8 时为纯文本)、404、409、422、500 |
POST |
/v1/reload |
router::reload |
200,ReloadReport |
422、500 |
GET |
/v1/connections |
tracking::connections |
200,{"sessions": [...], "flows": [...]} |
无 |
DELETE |
/v1/connections/{flow} |
tracking::close_one |
204,空 body |
404;flow id 不是 u64 时为 400(纯文本) |
DELETE |
/v1/connections?outbound=&inbound=&session=&user= |
tracking::close_matching |
200,{"closed": <n>} |
400 |
GET |
/v1/connections/events |
tracking::events |
200,由 opened、closed 和 lagged 事件组成的 text/event-stream |
无 |
GET |
/v1/traffic |
tracking::traffic |
200,由 traffic 事件组成的 text/event-stream |
无 |
GET |
/v1/version |
router::version |
200,{"version": "<the app's version>"} |
无 |
GET |
/v1/health |
router::health |
200,{"status": "ok"} |
无 |
设置了 secret 时,对以上任一路径或其他任何路径的请求,只要没有携带 bearer token,都会先得到 401,先于其他任何处理。
axum 的 get 也会在同一路径上响应 HEAD:它运行 GET 处理函数,但不发送 body。没有任何路由响应 OPTIONS,所以 OPTIONS 请求得到 405(设置了 secret 时则先得到 401)。
{group} 由 axum 的 Path 提取器做百分号解码,所以名字带空格的组要写成 /v1/routes/Video%20Streaming。DEFAULT_GROUP = "default" 指代默认路由,因此 PUT /v1/routes/default 编辑的是没有任何组匹配的流量所用的选择。[subscribe.routes] 用同一个键表示默认节点,而含有名为 default 的路由组的订阅文件无法合并(订阅文件)。
/v1/version 返回 app 二进制的 env!("CARGO_PKG_VERSION"),它作为 app_version 传给 router():在此版本中为 2.1.1。只要 API 在提供服务,/v1/health 就返回 {"status": "ok"};它不检查 supervisor 的状态。
状态码与错误 body
Section titled “状态码与错误 body”API 自身代码产生的错误都是同一形状的 JSON,由 router::error 构造:
{"error": "<message>"}| 状态码 | 何时 | Body |
|---|---|---|
400 |
PUT /v1/routes/{group}:body 不是 {"target": "<string>"},或 API 不接受的 If-Match;DELETE /v1/connections:查询串无法反序列化 |
JSON:axum 的拒绝文本,或 If-Match 的用法说明文本 |
400 |
DELETE /v1/connections/{flow},flow id 不是 u64 |
axum 给出的纯文本:Invalid URL: Cannot parse `abc` to a `u64` |
400 |
PUT /v1/routes/{group},{group} 百分号解码后不是 UTF-8,例如 %FF。router::set 接受 Path<String> 而不捕获它的拒绝,所以 axum 在处理函数运行之前就作答 |
axum 给出的纯文本:Invalid URL: Invalid UTF-8 in `group` |
401 |
设置了 secret,而请求没有携带它 | JSON a valid bearer token is required,以及 WWW-Authenticate: Bearer |
404 |
ControlError::UnknownGroup;或对一条已不存活的流执行 DELETE /v1/connections/{flow} |
JSON |
404 |
router 没有的路径(未设置 secret 时) | 空:axum 的默认 fallback |
405 |
已知路径配上其他方法,例如 DELETE /v1/connections/events,或任何 OPTIONS(未设置 secret 时) |
空,带 axum 的 Allow 头 |
409 |
ControlError::Stale |
JSON,以及文件当前的 ETag 头 |
422 |
ControlError::Invalid |
JSON |
500 |
ControlError::Failed |
JSON |
400 的几种情况使用 axum 自己的拒绝文本:router::set 和 tracking::close_matching 捕获这些拒绝,并以 JSON 重新包装,状态码一律为 400,不管 axum 本来会给什么状态码:
| 拒绝原因 | 文本开头 | axum 自己的状态码 |
|---|---|---|
没有 Content-Type: application/json |
Expected request with `Content-Type: application/json` |
415 |
| body 不是 JSON | Failed to parse the request body as JSON: |
400 |
JSON 形状不对:没有 target、target 不是字符串,或有其他任何键(SetRoute 是 deny_unknown_fields) |
Failed to deserialize the JSON body into the target type: |
422 |
| body 超过 axum 的默认上限 2,097,152 字节 | Failed to buffer the request body: |
413 |
API 不认识的选择器字段、重复的字段,或不是 u64 的 session |
Failed to deserialize query string: |
400 |
router() 在合并入跟踪路由之后,用 middleware::from_fn_with_state 把 bearer 函数作为一层套在整个 router 外面。Router::layer 覆盖每条路由和 fallback,所以未知路径得到的是 401 而不是 404(a_secret_is_required_as_a_bearer_token)。
- 状态是
Arc<sha2::digest::Output<Sha256>>:secret.expose().as_bytes()的 SHA-256,在router()中计算一次。router 只保存这个摘要。 bearer读取Authorization,把它当作字符串(HeaderValue::to_str,接受可见 ASCII、空格和制表符;含其他任何字节的值不算 token),在第一个空格处拆分,接受任意大小写的 schemebearer,去掉 token 两端的空白,然后对它求哈希。- 两个摘要的比较方式是:把每对字节的 XOR 折叠进一个字节,再检查它是否为零。比较的是摘要而不是 token,所以无论出示的是什么(包括它的长度),比较耗时都相同。
- 匹配则继续执行后面的栈。其他任何情况,包括缺少这个头、使用其他 scheme 或 scheme 后面没有空格,都得到
401,body 为{"error": "a valid bearer token is required"},并带WWW-Authenticate: Bearer。
Secret<String> 在 Debug 下打印为 Secret(..),所以 ApiOptions 派生的 Debug 从不显示 token。未设置 secret 时,根本没有这个中间件。
router
Section titled “router”struct Shared<C> { control: C, /// Held across an edit and the reload after it. edits: tokio::sync::Mutex<()>, version: &'static str,}
pub fn router<C: ConfigControl>( control: C, tracker: Tracker, secret: Option<Secret<String>>, app_version: &'static str,) -> axum::Router;router() 构建两个 router 并把它们合并。配置路由(/v1/routes、/v1/routes/{group}、/v1/reload、/v1/version、/v1/health)以 Arc<Shared<C>> 作为状态;跟踪路由只以 Tracker 作为状态,因为它们不涉及配置。secret 为 Some 时,合并后的 router 会加上 bearer 层。它返回一个普通的 axum::Router,由调用方负责提供服务:app 把它交给 ApiListener::serve,测试则用 tower::ServiceExt::oneshot 驱动它。
Shared::edits 是一个 tokio::sync::Mutex<()>。router::set 在解析完 body 和 If-Match 之后获取它,并在 set_route 及随后的 reload 期间一直持有,因此 API 编辑一次只应用一个,每个响应报告的都是它自己的那次重载。POST /v1/reload 不获取这个锁,也没有任何机制对 API 编辑与来自别处的重载排序;重载本身由 Core::reload_with 串行化。
ConfigControl
Section titled “ConfigControl”pub trait ConfigControl: Send + Sync + 'static { fn routes(&self) -> Result<RouteView, ControlError>;
fn set_route( &self, group: &str, target: &str, if_match: Option<&ETag>, ) -> impl Future<Output = Result<ETag, ControlError>> + Send;
fn reload(&self) -> impl Future<Output = Result<ReloadReport, ControlError>> + Send;}| 方法 | 约定 |
|---|---|
routes |
磁盘上的配置文件所描述的路由,以及文件的 ETag。同步:ConfigFile 在请求的任务上读取并解析文件。 |
set_route |
在配置文件中把 group(或 DEFAULT_GROUP)指向 target,并返回文件新的 ETag。不重载。编辑后的文件在写入之前会像启动时那样被校验;校验失败则文件保持不动。带 if_match 时,如果文件的 ETag 已不再是它,文件同样保持不动,错误中给出当前的 ETag。它是异步的,因为校验要构建启动时会构建的一切。 |
reload |
从磁盘重载,并报告它做了什么。 |
这个 trait 使用返回位置的 impl Future,而不是 async_trait 的装箱,所以实现可以直接写 async fn(ConfigFile 和测试里的 NoConfig 都是这样)。Send + Sync + 'static 约束让 router() 能把它放进一个由所有请求共享的 Arc 中。
RouteView 及其组成部分
Section titled “RouteView 及其组成部分”pub const DEFAULT_GROUP: &str = "default";
pub struct RouteView { pub default: DefaultView, pub groups: Vec<GroupView>, pub targets: Vec<Target>, #[serde(skip)] pub etag: ETag,}
pub struct DefaultView { pub pick: Option<String>, pub target: Option<String>,}
pub struct GroupView { pub name: String, pub pick: Option<String>, pub inherit: Inherit, pub target: Option<String>,}
#[serde(rename_all = "snake_case")]pub enum Inherit { Default, Direct, Blackhole, Group(String) }
pub struct Target { pub name: String, pub kind: TargetKind,}
#[serde(rename_all = "snake_case")]pub enum TargetKind { Node, Outbound, Balancer, Direct, Blackhole }它们都派生了 Debug、Clone、PartialEq、Eq 和 Serialize(TargetKind 还派生了 Copy)。
pick是文件中写的选择。target是流量的去向:有选择时就是选择本身,否则是该路由回退到的去向。- 只有在文件完全没有指定任何出口,或者某个组循环继承、或继承自不存在的组时,
target才为null。这样的配置无法应用,但仍然会被描述出来,让用户能看到并重新选择。 Inherit按订阅文件的写法序列化:"default"、"direct"、"blackhole"或{"group": "<name>"}。TargetKind序列化为"node"(订阅文件的节点)、"outbound"(配置自身的出站)、"balancer"、"direct"或"blackhole"(内置出站)。etag作为ETag头发送,不在 body 里。
下面是 lists_each_groups_pick_and_inherit_and_the_targets 的响应。它的配置有一个 freedom 出站 local-direct;订阅文件有节点 us-socks 和 jp-socks,以及组 Video Streaming(继承 Proxy)和 Proxy;选择为 default = "us-socks" 和 "Video Streaming" = "jp-socks":
{ "default": { "pick": "us-socks", "target": "us-socks" }, "groups": [ { "name": "Video Streaming", "pick": "jp-socks", "inherit": { "group": "Proxy" }, "target": "jp-socks" }, { "name": "Proxy", "pick": null, "inherit": "default", "target": "us-socks" } ], "targets": [ { "name": "local-direct", "kind": "outbound" }, { "name": "us-socks", "kind": "node" }, { "name": "jp-socks", "kind": "node" }, { "name": "direct", "kind": "direct" }, { "name": "blackhole", "kind": "blackhole" } ]}ReloadReport
Section titled “ReloadReport”pub enum ReloadReport { Applied(ApplyReport), Unchanged,}Unchanged 表示这些文件与上次尝试的文件逐字节相同,因此什么也没做;重载如何做出这一判断见 etemenanki-app:运行、重载与关闭。ReloadReport 派生了 Debug、Clone、PartialEq 和 Eq。Serialize 是手写的:
Unchanged为{"status": "unchanged"}。Applied为{"status": "applied"},后面跟着 supervisor 的ApplyReport的七组内容,顺序为built、swapped、drained、rebound、restarted、removed、reused。每组是一个字符串列表,由私有的Names包装器把每一项的Display写出。
这些项在 built 和 reused 中是 Resource 值(inbound <tag>、outbound <tag>@v<n>、balancer <tag>、user set <tag>、route、dns),在 drained 中是 OutboundId 值(<tag>@v<n>),在 swapped、rebound、restarted 和 removed 中是裸的入站 tag(规划并应用变更)。计划按 DNS、出站、负载均衡器、路由、用户集、入站的顺序列出资源。在上面的示例中把 Proxy 组切换到 jp-socks,只会改变路由表。下面的报告是根据这个计划顺序推出的;测试只断言 "status": "applied":
{ "status": "applied", "built": ["route"], "swapped": [], "drained": [], "rebound": [], "restarted": [], "removed": [], "reused": ["dns", "outbound local-direct@v1", "outbound us-socks@v1", "outbound jp-socks@v1"]}目标不同,结果也不同。上面的示例中没有任何东西路由到 direct,所以不存在 direct 出站。把 Proxy 切换到 direct(comments_and_layout_survive_a_switch 就是这样做的)会让订阅合并创建内置的 freedom 出站(subscribe::ensure_builtins),因此 built 中还会列出 outbound direct@v1。
ControlError
Section titled “ControlError”#[derive(Debug, thiserror::Error)]pub enum ControlError { #[error("no route group named {0:?}")] UnknownGroup(String), #[error("{0}")] Invalid(String), #[error("the config file has changed since it was read; its ETag is now {current}")] Stale { current: ETag }, #[error("{0}")] Failed(String),}| 变体 | 状态码 | 含义 | 示例文本 |
|---|---|---|---|
UnknownGroup |
404 |
没有叫这个名字的路由组 | no route group named "Nope" |
Invalid |
422 |
该变更无法得到合法的配置(以这种方式被拒绝的变更不会改动文件),或者磁盘上的配置无法应用 | no node, outbound or balancer is named "nowhere"; the route view lists them |
Stale |
409 |
文件已不再是 If-Match 所指的那个,或者在检查编辑期间文件发生了变化 |
the config file has changed since it was read; its ETag is now "<64 hex digits>" |
Failed |
500 |
其他一切:无法读写的文件、无法绑定的监听器、已停止的运行时 | writing /etc/etemenanki/config.toml: Permission denied (os error 13) |
IntoResponse for ControlError 根据变体选择状态码,写出 {"error": "<Display>"},并且只对 Stale 额外加上 current 的 ETag 头,这样客户端无需再发一次 GET,就能针对文件的当前内容重试。
#[derive(Debug, Clone, PartialEq, Eq, Hash)]pub struct ETag(String);
impl ETag { pub fn of(bytes: &[u8]) -> Self; pub(crate) fn header(&self) -> HeaderValue; pub(crate) fn from_if_match(value: Option<&HeaderValue>) -> Result<Option<Self>, &'static str>;}
impl fmt::Display for ETag { /* "\"<hex>\"" */ }of是文件字节的 SHA-256 的小写十六进制:一个强实体标签(strong entity tag)。两次读取文件得到相同的标签,当且仅当读到的字节相同;因此,基于一个标签已不再是文件当前标签的视图所做的编辑,就是针对一个此后已经变化的文件所做的编辑。Display和header按 HTTP 书写实体标签的方式,把它写成带引号的形式:ETag: "9f86d081884c7d659a2feaa0c55ad015a3bf4f1b2b0b822cd15d6c15b0f00a08"。header预期转换一定成功(“a quoted hex string is a valid header value”);它只会对 API 自己计算出的标签调用。- 标签只覆盖配置文件。配置所指定的订阅文件不在其中,所以
If-Match捕捉不到GET与PUT之间对订阅文件的修改。Snapshot::pick依据set_route读到的订阅文件检查目标,而它之后的检查(check_bytes)会再次从磁盘读取订阅文件。
from_if_match 读取 If-Match 头:
| 头 | 结果 |
|---|---|
| 缺失 | Ok(None):无条件 |
*(去掉空白后,包括制表符) |
Ok(None):任何当前文件都可以 |
一个带引号、内部不含 " 或 , 的标签,例如 "9f86…0a08" |
Ok(Some(tag)),与文件的标签按字符串相等比较 |
弱标签(W/"…")、列表("a", "b")、不带引号的值,或含可见 ASCII、空格和制表符以外字节的值(HeaderValue::to_str 失败) |
Err,以 400 回答,文本为 If-Match takes "*" or the one entity tag an ETag header gave |
列表和弱标签被直接拒绝,而不是半支持:在 If-Match 所用的强比较下,弱标签永远不会匹配;而 API 只会发出一个强标签,也就是要回传的那个。格式正确但不是文件标签的值(大写十六进制、旧标签、"")不算解析错误;它会到达 set_route,并在那里以 409 失败。
Listen、ApiOptions 与 OptionsError
Section titled “Listen、ApiOptions 与 OptionsError”pub const DEFAULT_LISTEN: &str = "127.0.0.1:9090";
pub enum Listen { Tcp(SocketAddr), Unix(PathBuf),}
impl Listen { pub fn is_local(&self) -> bool;}impl FromStr for Listen { type Err = OptionsError; }impl fmt::Display for Listen { /* 地址,或路径 */ }
pub struct ApiOptions { pub listen: Listen, pub secret: Option<Secret<String>>,}
impl ApiOptions { pub fn new(listen: Listen, secret: Option<Secret<String>>) -> Result<Self, OptionsError>;}
pub enum OptionsError { #[error("listen {0:?}: expected ip:port, or a path (containing a /) for a unix socket")] Listen(String), #[error("listen {0}: an API reachable from other hosts needs a secret; listen on a loopback address or a unix socket, or set one")] SecretRequired(Listen), #[error("secret is empty")] EmptySecret,}Listen::from_str 把任何含 / 的值当作 unix socket 路径(/run/etemenanki/api.sock,或相对于工作目录的 ./api.sock),因为任何 ip:port 都不含 /。其他值必须能解析为 SocketAddr(127.0.0.1:9090、[::1]:9090)。主机名是错误,而不会被当成以它命名的 socket 文件;不带 / 的裸文件名同样是错误。
对于 IP 满足 is_loopback()(127.0.0.0/8 或 ::1)的 TCP 地址,以及所有 unix socket,is_local 为 true。0.0.0.0、[::] 和 IPv4 映射的 IPv6 地址都不算 loopback。
ApiOptions::new 按以下顺序检查:
- 空 secret 以
EmptySecret拒绝,因为那是任何人都能出示的 secret。 - 没有 secret 且
listen不是本地地址,以SecretRequired拒绝。
这两个类型都派生了 Debug、Clone、PartialEq 和 Eq。Secret 按其内容比较,所以两个 ApiOptions 即使只有 secret 不同也不相等;二进制程序依靠这一点在 secret 改变时重新绑定。字段是公开的,所以代码可以直接构造 ApiOptions,跳过这两项检查;app 只通过 api::options 中的 new 来构造它。
ApiListener
Section titled “ApiListener”pub struct ApiListener(Bound);
enum Bound { Tcp(TcpListener), Unix { listener: UnixListener, _file: SocketFile },}
impl ApiListener { pub async fn bind(listen: &Listen) -> io::Result<Self>; pub fn local_addr(&self) -> Option<SocketAddr>; pub async fn serve( self, app: Router, shutdown: impl Future<Output = ()> + Send + 'static, ) -> io::Result<()>;}
struct SocketFile { path: PathBuf, id: (u64, u64), // (st_dev, st_ino)}TCP 的 bind 就是 tokio::net::TcpListener::bind。对 unix 路径,它会:
- 用
symlink_metadata查看该路径。那里若是 socket 文件,就认为是崩溃的上次运行遗留的陈旧文件,将其删除。若是其他任何文件,则返回AlreadyExists错误<path> exists and is not a socket。什么都没有则没问题;其他任何错误都原样返回。 - 在该路径上绑定一个
tokio::net::UnixListener。 - 记录新文件的设备号和 inode。如果这次
symlink_metadata失败,就删除该文件并返回错误。
SocketFile 的 Drop 只在路径仍然具有该设备号和 inode 时才删除它,所以此后被别的东西放到那里的 socket 文件不会被动。
bind 不会创建缺失的父目录。is_local 把每个 unix socket 都视为本地,所以 unix socket 上的 API 不需要 secret。不设 secret 的 loopback TCP 地址可以被主机上的每个进程访问。
local_addr 返回绑定的 TCP 地址(所以 :0 端口在日志中记录为内核选定的端口);对 unix socket,或 socket 无法报告自己的地址时,返回 None。serve 运行 axum::serve(listener, app).with_graceful_shutdown(shutdown)。shutdown 完成时,axum 停止接受连接,drop 监听器,要求每条打开的连接完成当前响应后关闭,并等待各连接任务结束。在 axum 0.8.9 中,返回的 future 总是以 Ok(()) 结束:axum 在它的循环内部处理 accept 错误,io::Result 只是为了类型而存在。
读取路由:Snapshot
Section titled “读取路由:Snapshot”pub struct Snapshot { pub bytes: Vec<u8>, pub etag: ETag, pub cfg: Config, pub subscribe: Option<Subscribe>,}
impl Snapshot { pub fn of(sources: Sources, parsed: std::io::Result<Config>) -> Result<Self, ControlError>; pub fn view(&self) -> RouteView; pub fn pick(&self, group: &str, target: &str) -> Result<Option<Vec<u8>>, ControlError>;}Snapshot 是一次读取得到的文件:配置的字节、它们的 ETag、解析结果,以及解析后的订阅文件。of 把无法解析的配置和无法解析的订阅文件(subscribe::parse,其错误以 subscribe: 开头)映射为 ControlError::Invalid。它不把订阅文件合并进配置:视图描述的是文件的原样,包括那些无法应用的选择。ConfigFile 从磁盘构建快照;移动端库则从交给它的内容构建快照。
view(私有的 routes::view)构建 RouteView:
- 目标。 先是配置自身的出站(
outbound),然后是它的负载均衡器(balancer),再是订阅文件的[[proxy_server]]节点(node),各自按文件顺序。然后是direct和blackhole(direct、blackhole),列表中已有同名项时则不加。这与subscribe::ensure_builtins一致:它只在没有其他东西使用该 tag 时才创建它们。 - 默认。 配置有
[subscribe]且其中有该键时,选择是[subscribe.routes] default,否则是[route] default。目标是该选择;没有选择时是配置自身的第一个出站;再没有则是第一个节点;否则为null。 - 组。 只在配置有
[subscribe]且读到了订阅文件时才有。每个[[route_group]]对应一个GroupView,按文件顺序,也就是组的匹配顺序。pick是该组在[subscribe.routes]中的键。inherit复制RouteGroupInherit。target是以默认目标调用subscribe::resolve_group的结果:组的选择;没有选择则是它继承的去向,沿着{ group = "…" }链经过其他组的选择。循环继承或继承不存在的组会让resolve_group失败,视图显示null;缺少默认目标时也是如此。一个共享的resolvedmap 让每条链只需走一遍。
没有 [subscribe] 的配置没有组:它唯一的选择是 [route] default(without_a_subscribe_file_the_default_route_is_the_one_pick)。
sequenceDiagram participant C as 客户端 participant R as router set participant F as ConfigFile participant D as 磁盘上的文件 participant I as Instance C->>R: PUT /v1/routes/Proxy,带 target 和 If-Match R->>R: 解析 body 和 If-Match,失败为 400 R->>R: 锁住 edits R->>F: set_route F->>D: 读取配置和订阅文件 F->>F: If-Match 与 ETag 比较,不符为 409 F->>F: Snapshot pick,失败为 404 或 422 F->>F: 对编辑结果执行 instance check_bytes,失败为 422 或 500 F->>D: 再读一次配置,有变化则为 409 F->>D: write_atomically F-->>R: 新的 ETag R->>I: reload I-->>R: ReloadReport,或一个错误 R-->>C: 200 或该错误,带新的 ETag
ConfigFile::set_route 逐步说明
Section titled “ConfigFile::set_route 逐步说明”pub struct ConfigFile { instance: Arc<Instance>,}
impl ConfigFile { pub fn new(instance: Arc<Instance>) -> Self;}
impl ConfigControl for ConfigFile { /* routes, set_route, reload */ }-
读取。 私有的
read调用config::read_sources(instance.config_path())和Snapshot::of。读取错误会加上配置路径作为前缀(<path>: <error>):错误种类为InvalidInput时是Invalid(422),例如[subscribe] needs a path naming the subscribe file;其他情况是Failed(500),例如subscribe file <path>: No such file or directory (os error 2)。能读取但无法解析的配置或订阅文件由Snapshot::of报为Invalid,不带前缀。 -
If-Match。 标签不是快照的标签时,返回
Stale { current: snap.etag }:409,文件不动。 -
选择。
Snapshot::pick(group, target):- 视图中未列出的、
default以外的group是UnknownGroup(404)。没有订阅文件时视图不列出任何组,所以除default外的每个名字都是404; - 视图的
targets中未列出的target是Invalid:no node, outbound or balancer is named "<target>"; the route view lists them; - 否则由私有的
edit生成新的字节。新字节与旧字节相同时(该组已经选择了该目标),pick返回None,set_route直接返回快照的ETag,不做检查也不写入。
set_string(见下文)把新值写成基本的双引号字符串。因此,如果文件以其他方式写这个选择,例如字面量字符串'jp-socks',即使它已经指向该目标,生成的字节也会不同,这次编辑会像其他编辑一样经过检查、写入和重载。 - 视图中未列出的、
-
检查。
instance::check_bytes(edited)对这些字节运行的内容与--test完全相同:config::sources_from(它会再次从磁盘读取订阅文件)、config::effective(订阅合并)、instance::build(lowering 和api::options),以及etemenanki_supervisor::check。后者校验 spec(期望状态),并构建每个出站、路由、handler(处理器)和用户表,但不绑定任何监听器(校验与应用错误)。check针对空的运行状态做规划,并且准备时不做绑定,所以它从不返回ApplyError::Disruptive或ApplyError::Bind;只有写入之后的重载才可能返回它们。失败经由From<LoadError>变成ControlError(见按请求重载),文件保持不动。an_edit_that_would_not_start_is_422_and_leaves_the_file_untouched固定了一个例子:direct在那里是已知目标,但它指向的是一个负载均衡器,检查以subscribe: "direct" is routed to, but it names a balancer, not a "freedom" outbound拒绝它。 -
重读。 检查要构建启动时会构建的一切,需要一段时间,所以
set_route会再次读取配置文件,并与快照的字节比较。如果文件在此期间变了,比如检查期间保存了一次手工编辑,就返回Stale,带上现有内容的标签,不会覆盖它。这次比较只覆盖检查期间:在这次读取之后、第 6 步的 rename 之前保存的手工编辑会被替换掉。这次读取的 I/O 错误是Failed(500),只有错误文本本身,例如No such file or directory (os error 2),不带第 1 步的路径前缀。 -
写入。
write_atomically(path, &edited)。出错时为Failed:writing <path>: <error>。 -
返回
ETag::of(&edited)。
ConfigFile 中的文件 I/O 是阻塞的 std::fs,在请求的任务上运行。
私有的 edit(bytes, subscribe, group, target) 把配置解析为 toml_edit::DocumentMut(不是 UTF-8 或不是 TOML 时为 Invalid),并设置一个字符串:
| 配置 | 表 | 键 |
|---|---|---|
有 [subscribe] |
[subscribe.routes] |
组的名字,或 default |
没有 [subscribe] |
[route] |
default(唯一存在的组) |
合并订阅文件时,[subscribe.routes] default 会覆盖 [route] default,所以有订阅文件时,切换默认路由写的是前者,后者保持不变。
child(parent, key)返回parent[key],不存在时插入一个与父级同类的空表:在表下面是[table],在其他东西下面是内联表。父级不是类表结构时为Invalid:the parent of "<key>" is not a table。- 只包含子表的表是隐式的:它没有自己的表头,例如只作为
[[route.rule]]的父级而存在的[route]。edit对它调用set_implicit(false),代码注释给出的理由是:现在它含有一个键,需要自己的表头。在此版本锁定的toml_edit0.25.15 中,这个调用不会改变输出:含有键的隐式表无论如何都会带着[route]表头打印出来,所以这个调用只是让该表在文档中变为显式。 - 存放路由的表不是类表结构时为
Invalid:the table holding the routes is not a table。这两个防护针对的都是Config在解析时就已经会拒绝的形状。 set_string原地替换已有的值,并把它的decor(周围的空白和注释)复制到新值上,所以行尾注释得以保留。键不存在或存放的是表时,用toml_edit::value插入,放在所在表的末尾。
文件中的其他一切,包括注释、顺序和格式,都按原样重新输出。comments_and_layout_survive_a_switch 比较的是整个文件:
# Picks.[subscribe.routes]default = "us-socks" # the node for everything else"Video Streaming" = "jp-socks" # region-locked catalogue
# Trailing notes stay too.# Picks.[subscribe.routes]default = "us-socks" # the node for everything else"Video Streaming" = "us-socks" # region-locked catalogueProxy = "direct"
# Trailing notes stay too.选择内置出站只写入选择本身:不会为 direct 或 blackhole 写入任何出站,路由指向它们时由合并来创建(without_a_subscribe_file_the_default_route_is_the_one_pick 检查出站数量保持为 2)。
fn write_atomically(path: &Path, bytes: &[u8]) -> io::Result<()>;fn create_beside(dir: &Path, name: &str) -> io::Result<(File, PathBuf)>;write_atomically 整体替换文件,所以读者(监视器、重启、编辑器)永远看不到写了一半的文件:
fs::canonicalize(path)。符号链接会被跟随,所以被替换的是它指向的文件,而不是链接本身。没有父目录或文件名的路径报<path> is not a file。- 读取原文件的
fs::Permissions。 create_beside(dir, name):在同一目录下新建一个空文件,让下面的 rename 保持在同一个文件系统内。- 对新文件:
set_permissions为原文件的权限,write_all,sync_all。 fs::rename(tmp, path)覆盖原文件。- 第 4、5 步中任何一步失败,就删除临时文件并返回错误。
- 打开目录并对它
sync_all,让 rename 本身落盘。
create_beside 把文件命名为 .<name>.api-edit.<pid>.<n>,其中 <n> 来自进程范围的 AtomicU64(NEXT,从 0 开始,Relaxed)。它以 write(true)、create_new(true) 和 mode(0o600) 打开:
create_new即O_CREAT | O_EXCL,所以候选名字上已有的任何东西(崩溃遗留的文件,或为了重定向写入而放置的符号链接)都不会被打开,更不会被截断。遇到AlreadyExists就换下一个<n>;其他任何错误都直接返回。- 文件以
0o600模式创建为空文件,再由 umask 进一步收窄,因此在第 4 步赋予它原文件的模式之前,它的权限不会比仅所有者可访问更宽。新配置可能包含 secret,要等到那之后才写入。代码注释称它仅所有者可访问,“because it is about to hold a config”。 - 连续 64 个名字都已存在时,它放弃并返回
AlreadyExists:no free temporary name for <name> in <dir>。
替换后的文件:
- 在写入任何内容之前获得原文件的模式位(
fs::Permissions); - 替换目录项。别处指向旧文件的硬链接仍然指向旧内容。
编辑之后的重载
Section titled “编辑之后的重载”set_route 返回后,router::set 在仍持有 edits 的情况下调用 control.reload(),并以重载的结果作答。无论重载是否成功,响应都带有文件新的 ETag 头,因为文件无论如何都已写入:
| 重载结果 | 响应 |
|---|---|
Applied(report) |
200,报告,新的 ETag |
Unchanged |
200,{"status": "unchanged"},新的 ETag。磁盘上的文件就是上次应用的那些:例如选择原本就是该目标且其他什么都没变,或者另一次重载先应用了写入的文件。 |
| 错误 | 该错误的状态码和 {"error": …},新的 ETag。文件中已是新的选择;运行的仍是旧配置。 |
即使因为该组已经选择了该目标(第 3 步)而 set_route 什么都没写,router::set 也会调用 reload。此时的回答取决于磁盘上的文件:与上次尝试的文件不同时(例如一次尚未应用的手工编辑)是 applied;是上次应用的文件时是 unchanged;上次尝试这些相同文件时被拒绝,则是按请求重载中的 422。
ConfigFile::reload 就是 instance.reload().await?.report(),与 POST /v1/reload 是同一个调用。重载会重新读取文件,不复用检查过的字节,并且也会应用磁盘上的其他任何变化。切换之后紧接着的 POST /v1/reload 结果是 Unchanged,而不是第二次应用(a_reload_after_the_apis_own_finds_nothing_changed)。
如果选择是唯一的变化,并且它在配置或订阅文件定义的出站、节点或负载均衡器之间移动,重载只构建 route 并复用每一个出站(即上面的报告),所以已打开的流保留它们打开时所用的出站,而该组的新流则走新的目标(a_switch_sends_new_flows_of_the_group_through_the_new_target,以及 plane(数据平面):为每个流选路)。新指向内置的 direct 或 blackhole、或不再指向它们的选择,还会构建或移除那个出站。
使用 If-Match
Section titled “使用 If-Match”不想在自己没见过的文件上执行切换的客户端,可以回传它最近一次 GET /v1/routes 或 PUT 得到的 ETag:
- 标签匹配时,编辑得以通过;
- 如果文件此后变了(一次手工编辑、另一个客户端的切换),结果是
409,文件不动,响应的ETag头是当前标签,客户端可以在下一次尝试时带上它(an_edit_against_a_stale_etag_is_409); - 不带
If-Match,或带*时,只有第 5 步的重读保护手工编辑,而且只保护检查期间保存的那一次。
POST /v1/reload 调用 control.reload(),不获取 edits。在 app 中,这就是 Instance::reload(etemenanki-app:运行、重载与关闭),然后是 Reload::report:
Reload |
ReloadReport 或错误 |
|---|---|
Applied(report) |
ReloadReport::Applied(report) |
Unchanged { failed: None } |
ReloadReport::Unchanged |
Unchanged { failed: Some(reason) } |
ControlError::Invalid("<reason> (the files have not changed since this was found)") |
加载失败通过 app/src/instance.rs 中的 impl From<LoadError> for ControlError 变成 ControlError,它决定是不是配置本身有错:
LoadError |
ControlError |
状态码 |
|---|---|---|
Config(e),e.kind() 为 InvalidData 或 InvalidInput:无法解析的配置或订阅文件、无法合并的订阅文件、lowering 拒绝的值、[api] 错误 |
Invalid(e.to_string()) |
422 |
Config(e),其他任何种类:无法读取的文件 |
Failed(e.to_string()) |
500 |
Apply(ApplyError::Bind { .. })、Apply(ApplyError::Stopped) |
Failed |
500 |
Apply(e),其他所有 ApplyError:重复的 tag、未知的引用、为空或无法探测的负载均衡器、无效的资源、构建失败、中断性变更 |
Invalid |
422 |
这些文本就是 etemenanki-app --test 在 configuration invalid: 之后打印的文本,以及 supervisor 的 ApplyError 文本(校验与应用错误)。
读取失败到达客户端时带有 Instance::reload 加上的前缀 cannot read <config path>: <error>,并保留错误的种类。GET 和 PUT 则通过 ConfigFile::read 读取,它加的前缀是裸路径,所以同一个缺失项从两边看到的文本不同:
| 失败 | GET、编辑前的 PUT |
POST /v1/reload、PUT 之后的重载 |
|---|---|---|
[subscribe] 没有 path |
422,<path>: [subscribe] needs a path naming the subscribe file |
422,cannot read <path>: [subscribe] needs a path naming the subscribe file |
| 订阅文件缺失 | 500,<path>: subscribe file <sub>: No such file or directory (os error 2) |
500,cannot read <path>: subscribe file <sub>: No such file or directory (os error 2) |
API 请求的每次重载都和其他重载一样记录日志,target 为 etemenanki_app::instance:info 级别的 config reloaded: <summary>,或 error 级别的以下之一:reload: <error>(文件无法读取)、reload: <error>; keeping the running config,以及 reload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply it。完整列表见 etemenanki-app:运行、重载与关闭。
webclient/src/tracking.rs 中的五个跟踪处理函数只以 Tracker 作为状态。router() 从 app 拿到它,即 instance.tracker(),它是 supervisor 那个跟踪器的克隆,所以这些处理函数能看到每个入站和出站的每一条流,并且跨越重载。流如何登记、强制关闭(kill)、公布和采样,见流跟踪、统计与限速;本节讲 HTTP 层面的形状。
时间以 Unix 毫秒表示,由 unix_ms(started) 在响应时按 SystemTime::now() - started.elapsed() 计算:跟踪器保存的是单调的 Instant,所以墙上时钟的跳变会移动报告的开始时间。早于 epoch 的时间为 0;超过 u64::MAX 毫秒的时间会饱和。
GET /v1/connections
Section titled “GET /v1/connections”{ "sessions": [ { "id": 1, "inbound": "socks-in", "source": "127.0.0.1", "user": null, "up": 18, "down": 17, "started_at": 1790000000000 } ], "flows": [ { "id": 1, "session": 1, "network": "tcp", "inbound": "socks-in", "user": null, "source": "127.0.0.1", "host": "127.0.0.1", "port": 40000, "sniffed": null, "outbound": "direct", "outbound_version": 1, "rule": null, "plane_epoch": 1, "started_at": 1790000000000, "up": 5, "down": 5 } ]}connections 先读取 tracker.flows(),再读取 tracker.sessions():这是对两个结构的两次独立读取,不是一个快照。流的 session 可能不在 sessions 中(会话在两次读取之间结束了,或在两次读取之后才打开),列出的会话也可能没有任何流。
sessions 是按会话 id 排序的 Tracker::sessions()。每一项是一个 SessionView:
| 字段 | 来自 SessionStats |
含义 |
|---|---|---|
id |
id.get() |
会话 id |
inbound |
inbound |
接受它的入站的 tag |
source |
source |
客户端 IP,或 null |
user |
设置了 user 时为 user_label |
用户对外使用的标签,没有用户时为 null。会话的第一条流为它绑定了用户之后,会话才会给出它的用户。 |
up、down |
up、down |
来自和发往客户端的线上字节,包括协议开销(按用户的用量计费) |
started_at |
started |
它被接受的时间 |
flows 是按流 id 用 sort_unstable_by_key 排序的 Tracker::flows()。每一项是一个由 FlowEntry 构建的 FlowView:
| 字段 | 含义 |
|---|---|
id |
流 id:在 supervisor 的生命周期内唯一,并按登记顺序递增 |
session |
承载它的会话;不由任何会话承载的流(TUN 设备的流)为 null |
network |
由目标的 DialNetwork 得出的 "tcp" 或 "udp";Unknown 和 Unix 为 "unknown" |
inbound |
入站 tag |
user |
用户的标签,没有用户时为 null |
source |
源 IP,或 null |
host、port |
目标:按请求原样的域名,或 IP。对 UDP 子链路,是打开它的那个数据包的目标 |
sniffed |
嗅探恢复出的域名,或 null;UDP 子链路上始终为 null |
outbound、outbound_version |
它打开时所用出站的 tag 和版本。负载均衡器的流给出的是成员。supervisor 自己的 DNS 服务的 tag 为空字符串 ""。 |
rule |
匹配的路由规则的索引;默认路由以及由 DNS 服务应答的查询为 null |
plane_epoch |
为它选路的路由 plane 的 epoch(plane:为每个流选路) |
started_at |
它登记的时间,即拨号完成之后 |
up、down |
发往目标和从目标返回的有效载荷字节 |
这个示例的形状与 connections_lists_live_sessions_and_flows 的输出一致。其中的 id、端口和时间是虚构的(测试的回显服务器绑定端口 0);字节数则是测试固定的值:一条 SOCKS 连接回显 5 字节,在会话上计为上行 3 + 10 + 5 = 18 个线上字节(问候、请求、有效载荷),下行 2 + 10 + 5 = 17 个,在流上则是每个方向 5 个有效载荷字节。
DELETE /v1/connections/{flow}
Section titled “DELETE /v1/connections/{flow}”close_one 调用 Tracker::kill(FlowId::new(flow)):
true:204,空 body。只要该 id 仍在登记中,kill就返回true,包括已被强制关闭、但所属连接还没有 drop 它的流。false:404,{"error": "no live flow <flow>"}。
强制关闭只关闭那一条流。对 mux 子流,承载连接及其其他子流继续运行。对 UDP 子链路,UDP 关联及其其他子链路继续运行,之后路由到该出站的数据包会打开一条新的子链路:一条带新 id 的新流。各协议如何处理被强制关闭的流,见流跟踪、统计与限速中的表格。deleting_a_flow_closes_it_alone 固定了两点:另一条连接继续回显,关闭之后的第二次 DELETE 是 404。
带选择器的 DELETE /v1/connections
Section titled “带选择器的 DELETE /v1/connections”#[derive(Deserialize)]#[serde(deny_unknown_fields)]pub(crate) struct Selector { outbound: Option<String>, inbound: Option<String>, session: Option<u64>, user: Option<String>,}close_matching 调用 Tracker::kill_where(|f| selector.matches(f)),并以 200 和 {"closed": <n>} 作答,<n> 是这次调用强制关闭的流的数量。一条流匹配,当且仅当它匹配给出的每个字段:
| 字段 | 匹配这样的流 |
|---|---|
outbound |
出站 tag 等于它,不论版本。负载均衡器的 tag 什么都不匹配,因为流给出的是成员。 |
inbound |
入站 tag 等于它 |
session |
会话 id 等于它。没有会话的流永远不匹配。 |
user |
设置了用户且其标签等于它。匿名流永远不匹配。 |
没有任何字段时,每条存活的流都匹配("no selector closes everything")。空值也是值:?outbound= 只匹配空 tag。不属于这四个的字段是 400,因为 Selector 是 deny_unknown_fields;重复的字段、不是 u64 的 session 同样是 400。kill_where 跳过已被强制关闭的流,所以同一个选择器发两次,第二次什么也不关闭。
GET /v1/connections/events
Section titled “GET /v1/connections/events”events 把 tracker.subscribe() 包装进 tokio_stream::wrappers::BroadcastStream,并把每一项映射为一个 SSE 事件:
| 项 | 事件 | 数据 |
|---|---|---|
Ok(FlowEvent::Opened(flow)) |
opened |
FlowView,其计数在 SSE 流发送该事件时读取(新流通常为 0) |
Ok(FlowEvent::Closed(flow)) |
closed |
带最终计数的 FlowView |
Err(BroadcastStreamRecvError::Lagged(missed)) |
lagged |
{"missed": <n>} |
事件携带的是共享的 Arc<FlowEntry>,而不是副本:Tracker::open 用它登记的条目发送 Opened,流句柄的 Drop 在流离开注册表之后发送 Closed。events 在 SSE 流的 map 中构建每个 FlowView,也就是在为这个客户端序列化事件时构建,所以 up、down 和 started_at 在那时才读取。较晚读到 opened 事件的客户端,看到的是该流此后已传输的量;对于在此期间已关闭的流,看到的是它的最终计数。closed 始终携带最终计数,因为那时负责计数的句柄已经不在了。
SSE 流从订阅时开始:客户端可能看到某条流的 closed,却从未看到它的 opened。跟踪器的 broadcast 通道保留最近 EVENT_CAPACITY = 1024 个事件;落后更多的客户端会收到一个 lagged 事件,带有它错过的数量,然后从仍保留的最旧事件继续。流的打开和关闭从不等待 SSE 客户端。
event: openeddata: {"id":7,"session":3,"network":"tcp","inbound":"socks-in","user":null,…,"up":0,"down":0}
event: closeddata: {"id":7,"session":3,"network":"tcp","inbound":"socks-in","user":null,…,"up":7,"down":7}GET /v1/traffic
Section titled “GET /v1/traffic”traffic 把 tracker.stats()(采样器的 watch 接收端)包装进 WatchStream::from_changes,它跳过订阅时的当前快照,产出此后发布的每一个快照。每个快照都变成一个 traffic 事件:
struct Rates { up: u64, down: u64, up_rate: u64, down_rate: u64, flows: usize }
struct Traffic { tick: u64, interval_ms: u64, #[serde(flatten)] total: Rates, inbounds: BTreeMap<String, Rates>, outbounds: BTreeMap<String, Rates>,}{ "tick": 42, "interval_ms": 1000, "up": 1000, "down": 1000, "up_rate": 0, "down_rate": 0, "flows": 1, "inbounds": { "socks-in": { "up": 1000, "down": 1000, "up_rate": 0, "down_rate": 0, "flows": 1 } }, "outbounds": { "direct": { "up": 1000, "down": 1000, "up_rate": 0, "down_rate": 0, "flows": 1 } }}| 字段 | 来自 StatsSnapshot |
含义 |
|---|---|---|
tick |
tick |
supervisor 启动以来的 tick(采样周期)数 |
interval_ms |
interval 的毫秒数,饱和计算 |
速率所覆盖的时长,按距上一个 tick 的实际间隔测量 |
up、down、up_rate、down_rate、flows |
total,展平到顶层 |
启动以来的有效载荷字节、这个 tick 内的每秒字节数、tick 时存活的流数 |
inbounds、outbounds |
按 tag,已排序 | 每个自启动以来承载过流的入站和出站 tag 的同样五项数字 |
StatsSnapshot::users 被省略:按用户的数据只通过 Rust API 提供。etemenanki-app 的 supervisor 每 DEFAULT_SAMPLE_INTERVAL = 1 秒采样一次,所以客户端每秒收到一个事件。watch 通道只保留最新的快照,所以读取慢于 tick 的客户端会看到 tick 跳号,而不会积压。
事件流的帧格式
Section titled “事件流的帧格式”两个 SSE 流都是 axum::response::sse::Sse 响应,带 Content-Type: text/event-stream 和 Cache-Control: no-cache。每个事件是 event: <name> 加一行 JSON 的 data:,由 json_event 构建,它预期序列化总会成功(“these views always serialise”)。KeepAlive::default() 在 15 秒内没有事件时发送一行 : 注释。
两个流都不会在 API 运行期间自行结束:
/v1/connections/events在其连接存活期间永不结束。broadcast 发送端位于跟踪器的共享状态中,而 router 自己的状态持有跟踪器的一个克隆。/v1/traffic在采样器的watch发送端被 drop 时结束,这发生在 supervisor 关闭、采样器停止时。
两者都随各自的连接一起被 drop。
etemenanki-app 中的 API
Section titled “etemenanki-app 中的 API”从 [api] 到 ApiOptions
Section titled “从 [api] 到 ApiOptions”app/src/config.rs 声明了这个段:
#[derive(Deserialize, Clone, PartialEq)]#[serde(deny_unknown_fields)]pub struct ApiConfig { #[serde(default = "default_api_listen")] // etemenanki_webclient::DEFAULT_LISTEN pub listen: String, #[serde(default)] pub secret: Option<String>,}Config::api 是 Option<ApiConfig>:没有 [api] 段时 API 不监听,只写了 [api] 的配置在 127.0.0.1:9090 上不带 secret 监听。api::options(cfg) 用 Listen::from_str 解析 listen,用 Secret::new 包装 secret,调用 ApiOptions::new,并把 OptionsError 转成种类为 InvalidInput、前缀为 [api] 的 io::Error。instance::build 在 lower 旁边调用它,所以 --test、启动和每次重载都会像拒绝其他配置错误一样拒绝错误的 [api],而拒绝它的重载不会动正在运行的 API。下表给出 etemenanki-app --test 打印的文本。Configuration OK. 是标准输出上的普通一行;configuration invalid: <error> 是一行 ERROR 级别的 tracing 日志(时间戳、级别、target etemenanki_app),表中引用的是它的消息:
[api] |
打印 |
|---|---|
listen = "0.0.0.0:9090" |
configuration invalid: [api] listen 0.0.0.0:9090: an API reachable from other hosts needs a secret; listen on a loopback address or a unix socket, or set one |
listen = "[::]:9090" |
configuration invalid: [api] listen [::]:9090: an API reachable from other hosts needs a secret; listen on a loopback address or a unix socket, or set one |
secret = "",任意 listen |
configuration invalid: [api] secret is empty |
listen = "localhost:9090" |
configuration invalid: [api] listen "localhost:9090": expected ip:port, or a path (containing a /) for a unix socket |
listen = "api.sock" |
configuration invalid: [api] listen "api.sock": expected ip:port, or a path (containing a /) for a unix socket |
未知的键,例如 cors = true |
configuration invalid: TOML parse error at line … |
只写 [api]、listen = "[::1]:9090"、listen = "./api.sock",或带 secret 的 listen = "0.0.0.0:9090" |
Configuration OK. |
--test 检查 [api] 但不绑定任何东西,所以发现不了端口已被占用。
Core 把运行中配置的选项保存在一个 tokio::sync::watch::Sender<Option<ApiOptions>> 中,由 Core::start 用第一份配置的选项创建。应用成功后,Core::reload_with 用 send_if_modified 更新它,只在新选项与旧选项不同时才通知接收端。被拒绝的重载,或发现没有任何变化的重载,都不动它。Instance::api() 交出一个接收端(etemenanki-app:运行、重载与关闭)。
main 在 Instance::start 之后立即启动 API:
options = instance.api();first = options.borrow().clone()。- 为
Some(first)时,调用start_api(instance.clone(), &first)。失败时记录failed to start: API on <listen>: <error>,以Duration::ZERO关闭实例,进程以失败退出:启动时无法绑定的[api]会让启动失败,和入站一样。 - 创建
stop_apioneshot,并 spawnfollow_api(instance, options, running, api_stopped)。 - 启动文件监视器。
start_api 绑定 ApiListener,记录 API listening on <addr>(来自 local_addr 的 TCP 地址;为 None 时则是配置的 listen,即 unix 路径,或无法报告自身地址的 socket 的 TCP 地址),构建 router(ConfigFile::new(instance), instance.tracker(), options.secret.clone(), env!("CARGO_PKG_VERSION")),创建一个 oneshot,把它的接收端作为优雅关闭的 future,然后 spawn listener.serve(app, …)。它返回 Api { stop, task }。
没有 [api] 的配置除了它的入站之外不绑定任何东西(a_config_without_api_binds_nothing_new)。
跟随重载:follow_api
Section titled “跟随重载:follow_api”stateDiagram-v2 [*] --> Off: 启动时没有 api 段 [*] --> On: 启动,API 绑定成功 [*] --> Failed: 启动,API 无法绑定 Failed --> [*]: 以失败退出 On --> Stopping: api 段改变或被移除 Stopping --> On: 新选项绑定成功 Stopping --> Off: 段被移除,或新的绑定失败 Off --> On: api 段被加入或改变,并且绑定成功 Off --> Off: 新选项绑定失败 On --> [*]: 关闭,在 Api stop 之后 Off --> [*]: 关闭
follow_api 在一个 tokio::select! 上循环,等待 stopped oneshot 和 options.changed():
stopped触发,或watch发送端已不存在:退出循环,停止正在运行的 API,返回。- 选项改变:用
borrow_and_update取出它们,若有正在运行的 API 则先停止它(Api::stop,见下文),然后:- 为
Some(wanted)时,调用start_api。失败时记录API on <listen>: <error>; it stays off until [api] changes。 - 为
None时,记录API stopped: the config no longer has [api]。
- 为
这个 select! 没有偏向:停止和选项变化同时就绪时,哪个先被处理都有可能。无论哪种情况,循环结束时正在运行的 API 都已停止。
旧监听器在新监听器绑定之前关闭,所以同一地址可以被再次占用:secret 改变会重新绑定同一个端口(the_api_follows_the_api_section_across_reloads)。重启发生在 follow_api 中,而不是在发布新选项的那次重载中,所以通过 API 本身发起的重载,例如手工编辑 [api] 之后的一次 POST /v1/reload,会在 API 停止之前完成它的请求。
const SHUTDOWN_GRACE: Duration = Duration::from_secs(5);
struct Api { stop: oneshot::Sender<()>, task: JoinHandle<io::Result<()>>,}Api::stop 在 stop 上发送,这会完成优雅关闭的 future,然后在 tokio::time::timeout(SHUTDOWN_GRACE, …) 下等待 serve 任务:
| 结果 | 动作 |
|---|---|
任务返回 Ok(()) |
无 |
任务返回 io::Error |
error:API: <error>。axum 0.8.9 的 serve future 总是返回 Ok(())(见 ApiListener),所以这个分支只是为了类型而存在,在此版本中不会走到 |
| 任务 panic 或被取消 | error:API task failed: <JoinError> |
| 5 秒已过 | task.abort(),然后 warn:API: requests still in flight when it stopped were dropped |
关闭时(SIGINT 或 SIGTERM),main 记录 shutting down,在 stop_api 上发送,等待 follow_api(它以上述宽限期停止正在运行的 API),之后才调用 instance.shutdown(SHUTDOWN_GRACE)。API 先停止,然后代理的连接才获得它们的宽限期。
全部来自二进制程序,target 为 etemenanki_app:
| 级别 | 文本 | 何时 |
|---|---|---|
info |
API listening on <addr or path> |
start_api 绑定了监听器 |
error |
failed to start: API on <listen>: <error> |
第一份 [api] 无法绑定;进程退出 |
error |
API on <listen>: <error>; it stays off until [api] changes |
某次重载的新 [api] 无法绑定 |
info |
API stopped: the config no longer has [api] |
某次重载移除了 [api] |
error |
API: <error> |
serve 任务返回了错误;在 axum 0.8.9 下不可达 |
error |
API task failed: <error> |
serve 任务 panic |
warn |
API: requests still in flight when it stopped were dropped |
serve 任务未能在 SHUTDOWN_GRACE 内结束 |
info |
shutting down |
收到关闭信号 |
API 请求的重载与每次重载一样,记录在 etemenanki_app::instance 下(按请求重载)。webclient crate 本身不记录任何日志,没有启用 tracing feature 构建的 axum 也不记录任何拒绝。
| 不变量 | 由谁保证 | 由哪些测试固定 |
|---|---|---|
| 切换影响运行状态的唯一途径是配置文件:先编辑,再重载 | ConfigControl::set_route 从不接触 supervisor;router::set 调用 reload |
a_switch_edits_the_file_reloads_and_returns_the_new_etag、a_switch_sends_new_flows_of_the_group_through_the_new_target |
| 视图未列出的目标被拒绝,文件不动 | Snapshot::pick |
an_invalid_target_is_422_and_leaves_the_file_untouched |
| 无法启动的编辑被拒绝,文件不动 | 在 write_atomically 之前执行 instance::check_bytes |
an_edit_that_would_not_start_is_422_and_leaves_the_file_untouched |
未知的组是 404,文件不动 |
Snapshot::pick |
an_unknown_group_is_404 |
针对过时 If-Match 的编辑是 409,带当前的 ETag,文件不动;当前标签可以通过 |
set_route 的标签比较;IntoResponse 添加头 |
an_edit_against_a_stale_etag_is_409 |
| 在检查编辑期间发生变化的文件不会被覆盖 | check_bytes 之后的重读 |
无 |
切换前后,ETag 头都是文件字节的 SHA-256 |
ETag::of |
lists_each_groups_pick_and_inherit_and_the_targets、a_switch_edits_the_file_reloads_and_returns_the_new_etag |
| 只有选择改变;注释、顺序和格式保留 | toml_edit,set_string 保留 decor |
comments_and_layout_survive_a_switch |
没有 [subscribe] 时,选择写为 [route] default,规则保留 |
edit 选择 [route];toml_edit 在表含有键后打印其表头 |
without_a_subscribe_file_the_default_route_is_the_one_pick |
| 文件永远不会被看到写了一半 | 临时文件、sync_all、rename、目录 sync_all |
无 |
| 临时名字上已有的文件永远不会被打开 | create_new(O_EXCL) |
无 |
| API 编辑一次应用一个,各自报告自己的重载 | Shared::edits 在 set_route 和 reload 期间持有 |
无 |
重载 API 刚应用过的文件,结果为 unchanged |
Core::reload_with 比较 Sources |
a_reload_after_the_apis_own_finds_nothing_changed |
| 设置了 secret 时,没有 token 就什么都不回答,未知路径也一样 | 套在合并后 router 上的 bearer 层 | a_secret_is_required_as_a_bearer_token |
| 其他主机可访问的 API 有非空的 secret | ApiOptions::new |
an_api_open_to_other_hosts_needs_a_secret |
[api] 拒绝未知键和主机名;listen 默认为 loopback;路径即 unix socket |
deny_unknown_fields、Listen::from_str、default_api_listen |
api_options_come_from_the_api_section、an_api_open_to_other_hosts_needs_a_secret |
没有 [api] 的配置不绑定任何新东西 |
main 只对 Some 选项启动 API |
a_config_without_api_binds_nothing_new |
API 跨重载跟随 [api],secret 改变会重新绑定同一端口 |
send_if_modified;follow_api 先停止再启动 |
the_api_follows_the_api_section_across_reloads |
| unix socket 文件只在仍是本进程绑定的那个时才被删除 | SocketFile 的设备号和 inode 检查 |
无 |
| 跟踪端点不需要配置 | 它们只以 Tracker 作为状态 |
webclient/tests/tracking.rs 中的每个测试都以 NoConfig 运行,它的所有方法都会失败 |
| SSE 客户端永远不会拖慢数据平面 | 跟踪器一侧的有界 broadcast 通道和 watch 通道 |
这里没有;见流跟踪、统计与限速 |
失败路径与取消
Section titled “失败路径与取消”| 失败 | 结果 | 文件 |
|---|---|---|
| 配置或订阅文件无法读取 | GET 和 PUT:500,<path>: <error>。POST /v1/reload 和 PUT 之后的重载:500,cannot read <path>: <error> |
不动 |
配置或订阅文件无法解析;[subscribe] 没有 path |
GET 和 PUT:422。两种重载:422,读取错误带前缀 cannot read <path>: |
不动 |
{group} 百分号解码后不是 UTF-8 |
400,axum 的纯文本,在处理函数运行之前 |
不动 |
body 或 If-Match 有误 |
400,在获取 edits 锁之前 |
不动 |
| 未知的组;未知的目标 | 404;422 |
不动 |
过时的 If-Match;检查期间文件变化 |
409,带当前 ETag |
不动 |
| 编辑后的配置未通过检查 | 422(其他种类的错误则为 500,例如 lowering 无法读取的文件) |
不动 |
| 检查后的重读失败 | 500,裸 I/O 错误,没有路径前缀 |
不动 |
配置路径无法规范化(canonicalize)、其元数据无法读取,或无法创建临时文件 |
500,writing <path>: <error>;什么都没有创建 |
不动 |
| 临时文件的权限设置、写入、sync 或 rename 失败 | 500,writing <path>: <error>;临时文件被删除 |
不动 |
| rename 之后的目录 sync 失败 | 500,writing <path>: <error> |
已替换 |
| 写入成功之后的重载失败 | 该重载的错误状态码,带新的 ETag |
已替换;运行的是旧配置 |
第一份 [api] 无法绑定 |
failed to start: API on …,以失败退出 |
— |
某次重载的新 [api] 无法绑定 |
记录日志;API 保持关闭,直到 [api] 改变 |
— |
取消:
-
PUT处理函数在请求中途被 drop。 处理函数是由其连接任务运行的 future。当客户端在得到回答之前关闭连接(hyper 在这个正忙的连接上读到流的结束,并结束该连接),或者main返回后运行时关闭时,它会被 drop。drop 它会释放edits。此时已经发生了什么,取决于它当时在等待什么:正在等待 文件状态 运行的内容 edits.lock()不动;什么都还没读 不变 set_route内部的instance::check_bytes不动;检查的工作被丢弃 不变 control.reload()已写入: write_atomically是同步的,中间没有 await,所以一旦开始就会执行到底重载停在它被等待的地方;supervisor 的 actor 已经开始的应用不会因此停止(etemenanki-app:运行、重载与关闭) -
SSE 客户端断开。 SSE 流被 drop,broadcast 接收端或
watch接收端也随之被 drop;跟踪器一侧没有任何东西在等它。 -
停止 API。 优雅关闭立即停止接受连接并释放监听 socket,然后等待连接结束;
Api::stop把这个等待限制在SHUTDOWN_GRACE内,之后中止 serve 任务。 -
停止进程。 API 先停止,然后
Instance::shutdown(SHUTDOWN_GRACE)给连接宽限期(etemenanki-app:运行、重载与关闭)。
| 项目 | 值 | 定义位置 | 含义 |
|---|---|---|---|
DEFAULT_LISTEN |
127.0.0.1:9090 |
webclient/src/listen.rs |
[api] 未设置 listen 时的值 |
DEFAULT_GROUP |
"default" |
webclient/src/lib.rs |
代表默认路由的组名 |
SHUTDOWN_GRACE |
5 秒 | app/src/main.rs |
Api::stop 等待 serve 任务的时长,也是退出时给 Instance::shutdown 的宽限期 |
| 临时名字 | 64 次尝试 | create_beside |
报 no free temporary name for … 之前尝试的名字数 |
| 临时文件模式 | 0o600,由 umask 收窄 |
create_beside |
适用于文件为空期间,直到复制原文件的模式为止 |
| 请求 body | 2,097,152 字节 | axum 的默认 body 上限 | 更大的 PUT body 为 400 |
| HTTP 版本 | HTTP/1 | workspace 的 axum feature(只有 http1) |
不支持 HTTP/2 |
| SSE keep-alive | 15 秒没有事件后发送一行 : 注释 |
axum 的 KeepAlive::default() |
空闲的流上传输的内容 |
EVENT_CAPACITY |
1024 | supervisor/src/track/mod.rs |
/v1/connections/events 客户端在收到 lagged 之前可以落后的事件数 |
| 采样周期 | 1 秒 | DEFAULT_SAMPLE_INTERVAL,etemenanki-app 沿用 |
每秒一个 /v1/traffic 事件 |
| ETag | 64 个小写十六进制数字,带引号 | etag.rs |
仅为配置文件的 SHA-256 |
app 一侧的单元测试在 app/tests/unit/api.rs 中(即 app/src/api.rs 的 tests 模块)。每个异步测试都在各自临时目录(etemenanki-api-<pid>-<n>)中的配置上启动一个真实的、没有入站的 Instance,以版本 "test" 在 ConfigFile 和 instance.tracker() 之上构建 router,并用 oneshot 驱动它。两个 [api] 测试是同步的:它们对解析出的段(api_section)调用 api::options,没有实例,也没有 router。夹具是一份客户端配置,在编辑可能弄丢注释的每个地方都写了注释;以及一个订阅文件,含节点 us-socks 和 jp-socks,以及组 Video Streaming(继承 Proxy)和 Proxy。
| 测试 | 固定的行为 |
|---|---|
lists_each_groups_pick_and_inherit_and_the_targets |
整个 RouteView JSON:选择、按订阅文件写法给出的 inherit、沿链解析出的目标、目标的顺序和种类;ETag 头就是文件的标签 |
a_switch_edits_the_file_reloads_and_returns_the_new_etag |
200,"status": "applied",新的 ETag 等于文件的标签;文件中有该选择;随后的 GET 以相同标签显示它 |
the_default_route_is_switched_through_its_own_name |
PUT /v1/routes/default 写入 [subscribe.routes] default |
an_invalid_target_is_422_and_leaves_the_file_untouched |
422,指出该目标;字节不变 |
an_edit_that_would_not_start_is_422_and_leaves_the_file_untouched |
一个实为负载均衡器的已知目标(direct)未通过检查:422,含 names a balancer;字节不变 |
an_unknown_group_is_404 |
404;字节不变 |
an_edit_against_a_stale_etag_is_409 |
手工编辑之后,旧标签得到 409,带当前标签,手工编辑完好无损;当前标签可以通过 |
a_secret_is_required_as_a_bearer_token |
没有头是 401,带 WWW-Authenticate: Bearer;少一个字符的 token 是 401;没有 token 的未知路径是 401;正确的 token 是 200 |
comments_and_layout_survive_a_switch |
替换选择时保留它的行尾注释;新的选择加入它所在的表;文件其余部分逐字节相同 |
a_reload_after_the_apis_own_finds_nothing_changed |
切换之后,POST /v1/reload 为 {"status": "unchanged"} |
without_a_subscribe_file_the_default_route_is_the_one_pick |
没有组;默认目标是第一个出站;切换写入 [route] default 并保留规则;选择 direct 不增加出站;其他任何组都是 404 |
api_options_come_from_the_api_section |
没有 [api] 为 None;只写 [api] 为不带 secret 的 DEFAULT_LISTEN;绝对路径和相对路径都是 unix socket;带 secret 的非本地 listen 被接受 |
an_api_open_to_other_hosts_needs_a_secret |
不带 secret 的 0.0.0.0 以 needs a secret 失败;空 secret 失败;主机名失败;未知键失败 |
app/tests/integration/e2e_api.rs 中的端到端测试运行 app 二进制,并直接用原始 HTTP/1.1 与它通信:
| 测试 | 固定的行为 |
|---|---|
a_config_without_api_binds_nothing_new(仅 Linux) |
从 /proc/<pid>/fd 和 /proc/<pid>/net/{tcp,tcp6,unix} 看:没有 [api] 时,进程只监听它的 SOCKS 端口,没有 unix socket;有 [api] 时,恰好监听那个端口和 API 的端口 |
a_switch_sends_new_flows_of_the_group_through_the_new_target |
继承 blackhole 的组丢弃 loopback 流量;PUT 到一个节点(第二个 app)之后,新连接能到达回显服务器,节点停止后则到达不了;PUT 到 direct 之后又能到达;GET 和文件都显示该选择 |
the_api_follows_the_api_section_across_reloads |
/v1/health 有回答;移除 [api] 会停止它;在另一个端口上带 secret 重新加回会启动它(没有 token 为 401);只改 secret 会重新绑定同一端口,旧 secret 为 401 |
webclient/tests/tracking.rs 中的跟踪测试运行一个真实的 supervisor,它有一个 SOCKS 入站,把一切都发往 freedom 出站 direct,并在 NoConfig 之上构建 router:
| 测试 | 固定的行为 |
|---|---|
connections_lists_live_sessions_and_flows |
一个会话和一条流,字段如上;会话线上字节 18 和 17,流有效载荷 5 和 5;rule 为 null;started_at 已设置 |
deleting_a_flow_closes_it_alone |
204;那条连接关闭;另一条继续回显;第二次 DELETE 为 404 |
deleting_by_selector_closes_the_matching_flows |
不匹配的 outbound 关闭 0 条;未知字段为 400;outbound 加 inbound 关闭两条流;没有选择器则关闭一切 |
traffic_streams_one_event_per_tick |
间隔为 100 ms 时,text/event-stream;tick 值连续;五个 tick 耗时在 4.5 到 10 个间隔之间;发送的 1000 字节出现在 up 中 |
connection_events_stream_opens_and_closes |
先是给出出站的 opened,然后是同一 id、最终计数为 7 和 7 的 closed |
ffi/tests/proxy.rs → set_route_and_reload_switch_like_the_rest_api 通过移动端库测试 routes::Snapshot(移动端库(etemenanki-ffi))。
没有测试覆盖:If-Match 解析(*、弱标签、列表、400);检查后的重读;write_atomically 和 create_beside(符号链接、权限、临时名字);lagged 事件;unix socket 上的 ApiListener 和 SocketFile;body 格式错误或 {group} 不是 UTF-8 时的 400;选择已经是该目标、或以字面量字符串写出时的 PUT;/v1/version;写入之后失败的重载;以及 Api::stop 的超时。改动这些路径时请补上测试。
运行它们:
cargo test -p etemenanki-app --lib api::testscargo test -p etemenanki-app --test integration e2e_apicargo test -p etemenanki-webclient测试布局和测试框架见测试。