The REST API (etemenanki-webclient)
Source files: 30 · checked against 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 is the HTTP API etemenanki-app serves when its config has an [api] section. It is meant for client configs, but main serves it for any config that has the section, whatever its inbounds. It lists the route groups and switches them, reloads on request, lists live sessions and flows, closes flows, and streams flow events and traffic rates as server-sent events. Route switching works by editing the config file and reloading it, so the picks live in the file alone: a restart reads the picks the API wrote.
The crate owns only HTTP. What a route group is, how the file is edited and validated, and what a reload does are the app’s, behind the ConfigControl trait that app/src/api.rs implements; the view and the edit themselves are in app/src/routes.rs, shared with the mobile library. The binary (app/src/main.rs) binds the API, and starts, stops or rebinds it when a reload changes [api]. This page is for contributors who change any of those files. How to configure and call the API is in the user guide. The tracker the connection endpoints read is on Tracking flows, stats and speed limits, and the reload an edit triggers is on etemenanki-app: running, reloading and shutting down and Planning and applying a change.
Responsibilities
Section titled “Responsibilities”| File | Owns |
|---|---|
webclient/src/lib.rs |
The public types: ConfigControl, RouteView and its parts, ReloadReport and its JSON, ControlError, DEFAULT_GROUP; the re-exports of the other modules |
webclient/src/router.rs |
router(): the routes, the bearer check, the edits mutex, and the mapping of ControlError to statuses |
webclient/src/etag.rs |
ETag: the SHA-256 of the config file’s bytes, and the parsing of If-Match |
webclient/src/listen.rs |
Listen, ApiOptions, OptionsError, DEFAULT_LISTEN, and ApiListener, which binds TCP or a unix socket and serves with a graceful shutdown |
webclient/src/tracking.rs |
The connection and traffic handlers and their JSON views, read straight from the supervisor’s Tracker |
app/src/api.rs |
options(), which turns [api] into ApiOptions; ConfigFile, the ConfigControl over a running Instance; write_atomically and create_beside |
app/src/routes.rs |
Snapshot: the route view of a config and its subscribe file, and the edit of one pick |
app/src/main.rs |
start_api, Api::stop, follow_api, and where the API sits in startup and shutdown |
The webclient crate does not:
- know the config schema. It never parses TOML; it calls
ConfigControl. - serve a UI. It answers JSON, server-sent events, the empty bodies of
204and of axum’s own404and405, and axum’s plain-text400for a path parameter it cannot parse. - send CORS headers. The router adds no CORS layer.
- serve per-user usage. The traffic view leaves the sampler’s per-user figures out, and per-user usage is the supervisor’s Rust API alone (Per-user usage accounting).
- watch files. The file watcher belongs to the binary (etemenanki-app: running, reloading and shutting down); the API only reloads when asked.
- log. It declares
tracingas a dependency but never calls it, and axum is built without itstracingfeature, so a rejected request is not logged either.
Dependency direction
Section titled “Dependency direction”flowchart LR app["etemenanki-app"] -->|"implements ConfigControl, binds ApiListener"| web["etemenanki-webclient"] ffi["etemenanki-ffi"] -->|"RouteView, ReloadReport, ControlError"| web ffi -->|"routes::Snapshot, instance::Core"| app web -->|"Tracker and its types, Secret, ApplyReport"| sup["etemenanki-supervisor"] web -->|"DialNetwork, Remote"| con["etemenanki-concepts"] app --> sup
The app depends on the webclient crate, not the other way round: ConfigControl is the seam. The crate uses these types from other workspace crates:
| Crate | Types | For |
|---|---|---|
etemenanki-supervisor |
topology::spec_plan::Secret |
The bearer token |
etemenanki-supervisor |
build::apply::ApplyReport |
What a reload did |
etemenanki-supervisor |
track::{Tracker, FlowEntry, FlowEvent, StatsSnapshot, TagStats}, entity::id::FlowId, entity::session::SessionStats |
Connections, flow events and traffic (tracking.rs) |
etemenanki-concepts |
net::{DialNetwork, Remote} |
A flow’s destination |
The mobile library reuses RouteView, ReloadReport, ControlError and routes::Snapshot for its native routes, reload and set_route calls, and does not serve [api]: when the config has the section, its lowering logs the config's [api] is not served here: the app reads routes and traffic through this library at info and publishes no API options. It never calls api::options, so it does not check the section either: an [api] that etemenanki-app would refuse is not refused there (Mobile library (etemenanki-ffi)). katana has no REST API.
The crate
Section titled “The crate”webclient/Cargo.toml declares etemenanki-webclient version 0.1.0, published to a private Cargo registry. It depends on etemenanki-supervisor (requirement 0.3.0), etemenanki-concepts (2.0.0), axum, tokio, tokio-stream (the SSE streams over the tracker’s channels), futures, compact_str, serde, sha2, hex, thiserror and tracing (declared, but unused). The tests add etemenanki-protocols, tower, http-body-util and serde_json.
The workspace Cargo.toml builds axum 0.8.9 with default-features = false and the features http1, json, query and tokio. The API therefore speaks HTTP/1 only, with no HTTP/2. axum’s other default features, form, matched-path, original-uri, tower-log and tracing, are off; the app is the only other crate that depends on axum, and it adds none.
Endpoints
Section titled “Endpoints”| Method | Path | Handler | Success | Errors |
|---|---|---|---|---|
GET |
/v1/routes |
router::list |
200, the RouteView as JSON, the file’s ETag header |
422, 500 |
PUT |
/v1/routes/{group} |
router::set |
200, the ReloadReport, the file’s new ETag header |
400 (plain text for a {group} that is not UTF-8), 404, 409, 422, 500 |
POST |
/v1/reload |
router::reload |
200, the ReloadReport |
422, 500 |
GET |
/v1/connections |
tracking::connections |
200, {"sessions": [...], "flows": [...]} |
none |
DELETE |
/v1/connections/{flow} |
tracking::close_one |
204, empty body |
404; 400 (plain text) for a flow id that is not a u64 |
DELETE |
/v1/connections?outbound=&inbound=&session=&user= |
tracking::close_matching |
200, {"closed": <n>} |
400 |
GET |
/v1/connections/events |
tracking::events |
200, text/event-stream of opened, closed and lagged events |
none |
GET |
/v1/traffic |
tracking::traffic |
200, text/event-stream of traffic events |
none |
GET |
/v1/version |
router::version |
200, {"version": "<the app's version>"} |
none |
GET |
/v1/health |
router::health |
200, {"status": "ok"} |
none |
With a secret set, a request to any of these, or to any other path, that does not carry the bearer token gets 401 before anything else.
axum’s get also answers HEAD on the same path: it runs the GET handler and sends no body. No route answers OPTIONS, so an OPTIONS request gets 405 (or 401 first when a secret is set).
{group} is percent-decoded by axum’s Path extractor, so a group with a space is addressed as /v1/routes/Video%20Streaming. DEFAULT_GROUP = "default" names the default route, so PUT /v1/routes/default edits the pick for traffic no group matches. [subscribe.routes] uses the same key for the default node, and a subscribe file with a route group named default does not merge (Subscribe files).
/v1/version answers env!("CARGO_PKG_VERSION") of the app binary, passed to router() as app_version: 2.1.1 at this revision. /v1/health answers {"status": "ok"} whenever the API serves; it does not look at the supervisor.
Statuses and error bodies
Section titled “Statuses and error bodies”Errors raised by the API’s own code are JSON of one shape, built by router::error:
{"error": "<message>"}| Status | When | Body |
|---|---|---|
400 |
PUT /v1/routes/{group}: a body that is not {"target": "<string>"}, or an If-Match the API does not take; DELETE /v1/connections: a query that does not deserialize |
JSON: axum’s rejection text, or the If-Match usage text |
400 |
DELETE /v1/connections/{flow} with a flow id that is not a u64 |
Plain text from axum: Invalid URL: Cannot parse `abc` to a `u64` |
400 |
PUT /v1/routes/{group} with a {group} whose percent-decoding is not UTF-8, such as %FF. router::set takes Path<String> without catching its rejection, so axum answers before the handler runs |
Plain text from axum: Invalid URL: Invalid UTF-8 in `group` |
401 |
A secret is set and the request does not carry it | JSON a valid bearer token is required, and WWW-Authenticate: Bearer |
404 |
ControlError::UnknownGroup; or DELETE /v1/connections/{flow} for a flow that is not live |
JSON |
404 |
A path the router does not have (without a secret) | Empty: axum’s default fallback |
405 |
A known path with another method, such as DELETE /v1/connections/events or any OPTIONS (without a secret) |
Empty, with axum’s Allow header |
409 |
ControlError::Stale |
JSON, and the file’s current ETag header |
422 |
ControlError::Invalid |
JSON |
500 |
ControlError::Failed |
JSON |
The 400 cases take axum’s own rejection texts, which router::set and tracking::close_matching catch and re-wrap as JSON with status 400 whatever status axum would have given them:
| Rejection | Text starts with | axum’s own status |
|---|---|---|
No Content-Type: application/json |
Expected request with `Content-Type: application/json` |
415 |
| Body is not JSON | Failed to parse the request body as JSON: |
400 |
JSON of the wrong shape: no target, a target that is not a string, or any other key (SetRoute is deny_unknown_fields) |
Failed to deserialize the JSON body into the target type: |
422 |
| Body over axum’s default limit of 2,097,152 bytes | Failed to buffer the request body: |
413 |
A selector field the API does not know, a repeated field, or a session that is not a u64 |
Failed to deserialize query string: |
400 |
Authentication
Section titled “Authentication”router() layers middleware::from_fn_with_state with the bearer function over the whole router, after the tracking routes are merged in. Router::layer covers every route and the fallback, so an unknown path gets 401, not 404 (a_secret_is_required_as_a_bearer_token).
- The state is
Arc<sha2::digest::Output<Sha256>>: the SHA-256 ofsecret.expose().as_bytes(), computed once inrouter(). The router keeps only the digest. bearerreadsAuthorization, takes it as a string (HeaderValue::to_str, which accepts visible ASCII, space and tab; a value with any other byte is no token), splits it at the first space, accepts the schemebearerin any case, trims the token of surrounding whitespace, and hashes it.- The two digests are compared by folding the XOR of every byte pair into one byte and testing it for zero. Digests are compared rather than tokens, so the comparison takes the same time whatever was presented, its length included.
- A match runs the rest of the stack. Anything else, including a missing header, another scheme or no space after the scheme, gets
401with{"error": "a valid bearer token is required"}andWWW-Authenticate: Bearer.
Secret<String> prints Secret(..) under Debug, so ApiOptions’ derived Debug never shows the token. When no secret is set, there is no middleware at all.
Key types
Section titled “Key types”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() builds two routers and merges them. The config routes (/v1/routes, /v1/routes/{group}, /v1/reload, /v1/version, /v1/health) take Arc<Shared<C>> as state; the tracking routes take the Tracker alone, since they involve no config. The merged router gets the bearer layer when secret is Some. It returns a plain axum::Router, which the caller serves; the app hands it to ApiListener::serve, and the tests drive it with tower::ServiceExt::oneshot.
Shared::edits is a tokio::sync::Mutex<()>. router::set takes it after the body and If-Match have been parsed, and holds it across set_route and the reload that follows, so API edits apply one at a time and each response reports its own reload. POST /v1/reload does not take it, and nothing orders an API edit against a reload from elsewhere; Core::reload_with serialises the reloads themselves.
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;}| Method | Contract |
|---|---|
routes |
The routing the config file on disk describes, with the file’s ETag. Synchronous: ConfigFile reads and parses the files on the request’s task. |
set_route |
Point group, or DEFAULT_GROUP, at target in the config file and return the file’s new ETag. Does not reload. The edited file is validated as a start would validate it before it is written; one that fails leaves the file untouched. With if_match, a file whose ETag is no longer that is left untouched too, and the error names the current one. Async because validating builds everything a start would. |
reload |
Reload from disk and report what that did. |
The trait uses return-position impl Future rather than an async_trait box, so an implementation writes async fn (both ConfigFile and the tests’ NoConfig do). The Send + Sync + 'static bound lets router() put it in an Arc shared by every request.
RouteView and its parts
Section titled “RouteView and its parts”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 }All of them derive Debug, Clone, PartialEq, Eq and Serialize (TargetKind is also Copy).
- A
pickis what the file says. Atargetis where the traffic goes: the pick when there is one, and otherwise what the route falls back to. - A
targetisnullonly where the file names no way out at all, or where a group inherits in a cycle or from a group that does not exist. Such a config does not apply, but it is still described, so the user can see it and pick again. Inheritserialises as the subscribe file writes it:"default","direct","blackhole", or{"group": "<name>"}.TargetKindserialises as"node"(a subscribe file’s node),"outbound"(an outbound of the config’s own),"balancer","direct"or"blackhole"(the built-ins).etagis sent as theETagheader, not in the body.
The response of lists_each_groups_pick_and_inherit_and_the_targets, for a config with one freedom outbound local-direct, a subscribe file with the nodes us-socks and jp-socks and the groups Video Streaming (inheriting Proxy) and Proxy, and the picks default = "us-socks" and "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 means the files are byte for byte the ones last tried, so nothing was done; how a reload decides that is on etemenanki-app: running, reloading and shutting down. ReloadReport derives Debug, Clone, PartialEq and Eq. Serialize is written by hand:
Unchangedis{"status": "unchanged"}.Appliedis{"status": "applied"}followed by the seven groups of the supervisor’sApplyReport, in the orderbuilt,swapped,drained,rebound,restarted,removed,reused. Each group is a list of strings, written by the privateNameswrapper as theDisplayof each item.
The items are Resource values in built and reused (inbound <tag>, outbound <tag>@v<n>, balancer <tag>, user set <tag>, route, dns), OutboundId values in drained (<tag>@v<n>), and bare inbound tags in swapped, rebound, restarted and removed (Planning and applying a change). The plan lists resources in the order DNS, outbounds, balancers, route, user sets, inbounds. A switch of the Proxy group to jp-socks in the fixture above changes only the route table. The report below is derived from that plan order; the tests assert only "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"]}The target matters. Nothing in that fixture routes to direct, so there is no direct outbound. A switch of Proxy to direct (as comments_and_layout_survive_a_switch makes) makes the subscribe merge create the built-in freedom outbound (subscribe::ensure_builtins), so built also lists 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),}| Variant | Status | Meaning | Example text |
|---|---|---|---|
UnknownGroup |
404 |
No route group has that name | no route group named "Nope" |
Invalid |
422 |
The change would not make a valid config (a change refused this way leaves the file untouched), or the config on disk cannot be applied | no node, outbound or balancer is named "nowhere"; the route view lists them |
Stale |
409 |
The file is no longer what If-Match named, or it changed while an edit was being checked |
the config file has changed since it was read; its ETag is now "<64 hex digits>" |
Failed |
500 |
Anything else: a file that cannot be read or written, a listener that cannot be bound, a runtime that has stopped | writing /etc/etemenanki/config.toml: Permission denied (os error 13) |
IntoResponse for ControlError picks the status from the variant, writes {"error": "<Display>"}, and for Stale alone adds the ETag header of current, so a client can retry against the file as it is now without another 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>\"" */ }ofis the lower-case hex of the SHA-256 of the file’s bytes: a strong entity tag. Two reads of the file have the same tag exactly when they read the same bytes, so an edit made against a view whose tag is no longer the file’s is an edit against a file that has changed since.Displayandheaderwrite it quoted, as HTTP writes an entity tag:ETag: "9f86d081884c7d659a2feaa0c55ad015a3bf4f1b2b0b822cd15d6c15b0f00a08".headerexpects the conversion to succeed (“a quoted hex string is a valid header value”); it is only called on tags the API computed.- The tag covers the config file only. The subscribe file it names is not part of it, so
If-Matchdoes not catch a change to the subscribe file between aGETand aPUT.Snapshot::pickchecks the target against the subscribe file asset_routeread it, and the check after it (check_bytes) reads the subscribe file from disk again.
from_if_match reads an If-Match header:
| Header | Result |
|---|---|
| Absent | Ok(None): no condition |
* (after trimming whitespace, tabs included) |
Ok(None): any current file |
One quoted tag with no " or , inside, such as "9f86…0a08" |
Ok(Some(tag)), compared by string equality with the file’s tag |
A weak tag (W/"…"), a list ("a", "b"), an unquoted value, or a byte other than visible ASCII, space and tab (HeaderValue::to_str fails) |
Err, answered 400 with If-Match takes "*" or the one entity tag an ETag header gave |
Lists and weak tags are refused rather than half-supported: a weak tag never matches under the strong comparison If-Match uses, and the API only ever hands out one strong tag, the one to send back. A well-formed tag that is not the file’s (upper-case hex, an old tag, "") is not a parse error; it reaches set_route and fails there with 409.
Listen, ApiOptions and OptionsError
Section titled “Listen, ApiOptions and 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 { /* the address, or the path */ }
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 treats any value containing a / as a unix socket path (/run/etemenanki/api.sock, or ./api.sock relative to the working directory), since no ip:port has one. Anything else must parse as a SocketAddr (127.0.0.1:9090, [::1]:9090). A host name is an error rather than a socket file named after it, and so is a bare file name without a /.
is_local is true for a TCP address whose IP is_loopback() (127.0.0.0/8 or ::1) and for every unix socket. 0.0.0.0, [::] and IPv4-mapped IPv6 addresses are not loopback.
ApiOptions::new checks, in this order:
- A secret that is empty is refused with
EmptySecret, since it would be one anybody could present. - No secret and a
listenthat is not local is refused withSecretRequired.
Both types derive Debug, Clone, PartialEq and Eq. Secret compares as its contents, so two ApiOptions differ when only the secret does; the binary relies on that to rebind on a secret change. The fields are public, so code can build an ApiOptions directly and skip both checks; the app builds them only through new, in api::options.
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)}bind for TCP is tokio::net::TcpListener::bind. For a unix path it:
- looks at the path with
symlink_metadata. A socket file there is taken to be stale, left by a crashed run, and removed. Any other file is anAlreadyExistserror,<path> exists and is not a socket. Nothing there is fine; any other error is returned. - binds a
tokio::net::UnixListenerat the path. - records the new file’s device and inode. If that
symlink_metadatafails, it removes the file and returns the error.
SocketFile’s Drop removes the path only while it still has that device and inode, so a socket file that something else has put there since is left alone.
bind does not create a missing parent directory. is_local counts every unix socket as local, so a unix-socket API needs no secret. A loopback TCP address without a secret is reachable by every process on the host.
local_addr returns the bound TCP address (so a :0 port is logged as the port the kernel chose), and None for a unix socket or when the socket cannot report its address. serve runs axum::serve(listener, app).with_graceful_shutdown(shutdown). When shutdown completes, axum stops accepting, drops the listener, asks each open connection to finish its current response and close, and waits for the connection tasks. With axum 0.8.9 the returned future always resolves to Ok(()): axum handles accept errors inside its loop, and the io::Result is there for the type.
Reading the routing: Snapshot
Section titled “Reading the routing: 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>;}A Snapshot is the files as read once: the config’s bytes, their ETag, their parse, and the parsed subscribe file. of maps a config that does not parse and a subscribe file that does not parse (subscribe::parse, whose errors start with subscribe: ) to ControlError::Invalid. It does not merge the subscribe file into the config: the view describes the files as written, including picks that would not apply. ConfigFile builds snapshots from disk; the mobile library builds them from the contents it was handed.
view (the private routes::view) builds the RouteView:
- Targets. The config’s own outbounds (
outbound), then its balancers (balancer), then the subscribe file’s[[proxy_server]]nodes (node), each in file order. Thendirectandblackhole(direct,blackhole), each unless something already in the list has that name, matchingsubscribe::ensure_builtins, which creates them only when nothing else carries the tag. - Default. The pick is
[subscribe.routes] defaultwhen the config has a[subscribe]with that key, else[route] default. The target is the pick, else the first outbound of the config’s own, else the first node, elsenull. - Groups. Only when the config has a
[subscribe]and a subscribe file was read. OneGroupViewper[[route_group]], in file order, which is the order the groups are matched.pickis the group’s key in[subscribe.routes].inheritcopiesRouteGroupInherit.targetissubscribe::resolve_groupwith the default target: the group’s pick, else what it inherits, following{ group = "…" }chains through the other groups’ picks. A cycle or an inherit of a missing group makesresolve_groupfail, and the view showsnull; so does a missing default target. A sharedresolvedmap lets each chain be walked once.
A config without [subscribe] has no groups: its only pick is [route] default (without_a_subscribe_file_the_default_route_is_the_one_pick).
Switching a route
Section titled “Switching a route”sequenceDiagram participant C as client participant R as router set participant F as ConfigFile participant D as files on disk participant I as Instance C->>R: PUT /v1/routes/Proxy with target and If-Match R->>R: parse body and If-Match, or 400 R->>R: lock edits R->>F: set_route F->>D: read config and subscribe file F->>F: If-Match against the ETag, or 409 F->>F: Snapshot pick, or 404 or 422 F->>F: instance check_bytes of the edit, or 422 or 500 F->>D: read the config again, changed means 409 F->>D: write_atomically F-->>R: the new ETag R->>I: reload I-->>R: ReloadReport, or an error R-->>C: 200 or the error, with the new ETag
ConfigFile::set_route step by step
Section titled “ConfigFile::set_route step by step”pub struct ConfigFile { instance: Arc<Instance>,}
impl ConfigFile { pub fn new(instance: Arc<Instance>) -> Self;}
impl ConfigControl for ConfigFile { /* routes, set_route, reload */ }-
Read. The private
readcallsconfig::read_sources(instance.config_path())andSnapshot::of. A read error is prefixed with the config path (<path>: <error>):Invalid(422) when its kind isInvalidInput, such as[subscribe] needs a path naming the subscribe file;Failed(500) otherwise, such assubscribe file <path>: No such file or directory (os error 2). A config or subscribe file that reads but does not parse isInvalidfromSnapshot::of, without the prefix. -
If-Match. With a tag that is not the snapshot’s,
Stale { current: snap.etag }:409, file untouched. -
Pick.
Snapshot::pick(group, target):- a
groupother thandefaultthat the view does not list isUnknownGroup(404). Without a subscribe file the view lists no groups, so every name butdefaultis404; - a
targetthe view’stargetsdoes not list isInvalid:no node, outbound or balancer is named "<target>"; the route view lists them; - otherwise the private
editproduces the new bytes. When they equal the old ones (the group already picks that target),pickreturnsNone, andset_routereturns the snapshot’sETagwithout checking or writing anything.
set_string(below) writes the new value as a basic, double-quoted string. A pick the file writes another way, such as the literal string'jp-socks', therefore comes out as different bytes even when it already names the target, and the edit is checked, written and reloaded like any other. - a
-
Check.
instance::check_bytes(edited)runs exactly what--testruns on those bytes:config::sources_from(which reads the subscribe file from disk again),config::effective(the subscribe merge),instance::build(lowering andapi::options), andetemenanki_supervisor::check, which validates the spec and builds every outbound, route, handler and user table without binding a listener (Validation and apply errors).checkplans against an empty running state and prepares without binding, so it never returnsApplyError::DisruptiveorApplyError::Bind; only the reload after the write can. A failure becomes aControlErrorthroughFrom<LoadError>(see Reloading on request) and leaves the file untouched.an_edit_that_would_not_start_is_422_and_leaves_the_file_untouchedpins one:directis a known target there, but it names a balancer, and the check refuses it withsubscribe: "direct" is routed to, but it names a balancer, not a "freedom" outbound. -
Re-read. The check builds everything a start would and takes a while, so
set_routereads the config file again and compares the bytes with the snapshot’s. A file that changed meanwhile, a hand edit saved during the check for example, isStalewith the tag of what is there now, and is not written over. The comparison covers the check only: a hand edit saved after this read and before the rename in step 6 is replaced. An I/O error on this read isFailed(500) with the bare error text, such asNo such file or directory (os error 2), without the path prefix of step 1. -
Write.
write_atomically(path, &edited). An error isFailed:writing <path>: <error>. -
Return
ETag::of(&edited).
The file I/O in ConfigFile is blocking std::fs, run on the request’s task.
Which key is edited
Section titled “Which key is edited”The private edit(bytes, subscribe, group, target) parses the config as a toml_edit::DocumentMut (not UTF-8, or not TOML, is Invalid) and sets one string:
| Config | Table | Key |
|---|---|---|
Has [subscribe] |
[subscribe.routes] |
The group’s name, or default |
No [subscribe] |
[route] |
default (the only group there is) |
[subscribe.routes] default overrides [route] default when the subscribe file is merged, so with a subscribe file a default switch writes the former and leaves the latter alone.
child(parent, key)returnsparent[key], inserting an empty table of the parent’s kind when it is absent: a[table]under a table, an inline table under anything else. A parent that is not table-like isInvalid:the parent of "<key>" is not a table.- A table that only held subtables is implicit: it had no header of its own, for example a
[route]that exists only as the parent of[[route.rule]].editcallsset_implicit(false)on it, and the code comment gives the reason as the header the table needs now that it holds a key. With the pinnedtoml_edit0.25.15 the call changes nothing in the output: an implicit table that holds a key is printed with its[route]header either way, so the call only makes the table explicit in the document. - A routes table that is not table-like is
Invalid:the table holding the routes is not a table. Both of these guard shapes thatConfigwould already have refused to parse. set_stringreplaces an existing value in place and copies itsdecor(the whitespace and comments around it) onto the new value, so a comment at the end of the line survives. A key that is not there, or holds a table, is inserted withtoml_edit::value, at the end of its table.
Everything else in the file, including comments, ordering and formatting, is re-emitted as it was. comments_and_layout_survive_a_switch compares whole files:
# 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.Picking a built-in writes only the pick: no outbound is written for direct or blackhole, which the merge creates when a route names them (without_a_subscribe_file_the_default_route_is_the_one_pick checks the outbound count stays at 2).
Writing the file
Section titled “Writing the file”fn write_atomically(path: &Path, bytes: &[u8]) -> io::Result<()>;fn create_beside(dir: &Path, name: &str) -> io::Result<(File, PathBuf)>;write_atomically replaces the file whole, so a reader (the watcher, a restart, an editor) never sees it half-written:
fs::canonicalize(path). A symlink is followed, so the file it names is replaced rather than the link. A path with no parent or file name is<path> is not a file.- Read the original’s
fs::Permissions. create_beside(dir, name): a new, empty file in the same directory, so the rename below stays on one file system.- On the new file:
set_permissionsto the original’s,write_all,sync_all. fs::rename(tmp, path)over the original.- If any of steps 4 and 5 failed, remove the temporary file and return the error.
- Open the directory and
sync_allit, so the rename itself is on disk.
create_beside names the file .<name>.api-edit.<pid>.<n>, where <n> comes from a process-wide AtomicU64 (NEXT, starting at 0, Relaxed). It opens with write(true), create_new(true) and mode(0o600):
create_newisO_CREAT | O_EXCL, so whatever already sits at a candidate name (a file left by a crash, or a symlink planted to redirect the write) is never opened, let alone truncated.AlreadyExistsmoves on to the next<n>; any other error is returned.- The file is created empty with mode
0o600, narrowed further by the umask, so it is never more open than owner-only before step 4 gives it the original’s mode. The new config, which may hold secrets, is written only after that. The code comment calls it owner-only “because it is about to hold a config”. - After 64 names in a row that exist, it gives up with
AlreadyExists:no free temporary name for <name> in <dir>.
The replacement:
- gets the original’s mode bits (
fs::Permissions) before anything is written to it; - replaces the directory entry. A hard link to the old file elsewhere keeps pointing at the old contents.
The reload after the edit
Section titled “The reload after the edit”After set_route returns, router::set calls control.reload() while still holding edits, and answers with the reload’s result. The response carries the file’s new ETag header whether or not the reload succeeded, since the file was written either way:
| Reload result | Response |
|---|---|
Applied(report) |
200, the report, new ETag |
Unchanged |
200, {"status": "unchanged"}, new ETag. The files on disk are those last applied: for example, the pick was already the target and nothing else changed, or another reload applied the written file first. |
| An error | Its status and {"error": …}, new ETag. The file holds the new pick; what runs is still the old config. |
router::set calls reload even when set_route wrote nothing because the group already picked the target (step 3). The answer is then whatever the files on disk give: applied when they differ from the last ones tried (a hand edit not yet applied, for example), unchanged when they are those last applied, and the 422 of Reloading on request when the last try of those same files was refused.
ConfigFile::reload is instance.reload().await?.report(), the same call as POST /v1/reload. The reload reads the files again and does not reuse the checked bytes, and it applies whatever else changed on disk too. A POST /v1/reload right after a switch is Unchanged, not a second apply (a_reload_after_the_apis_own_finds_nothing_changed).
When the pick is the only change and it moves between outbounds, nodes or balancers the config or the subscribe file defines, the reload builds only route and reuses every outbound (the report above), so open flows keep the outbound they were opened on, and new flows of the group follow the new target (a_switch_sends_new_flows_of_the_group_through_the_new_target, and The plane: routing each flow). A pick that newly names the built-in direct or blackhole, or stops naming it, also builds or removes that outbound.
Using If-Match
Section titled “Using If-Match”A client that wants to avoid switching against a file it has not seen sends back the ETag from its last GET /v1/routes or PUT:
- a matching tag lets the edit through;
- a file changed since (a hand edit, another client’s switch) is
409, the file untouched, and the response’sETagheader is the current tag, which the client can send with its next attempt (an_edit_against_a_stale_etag_is_409); - without
If-Match, or with*, only the re-read of step 5 protects a hand edit, and only one saved while the check ran.
Reloading on request
Section titled “Reloading on request”POST /v1/reload calls control.reload() without taking edits. In the app that is Instance::reload (etemenanki-app: running, reloading and shutting down), then Reload::report:
Reload |
ReloadReport or error |
|---|---|
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)") |
A failed load becomes a ControlError through impl From<LoadError> for ControlError in app/src/instance.rs, which decides whether the config is at fault:
LoadError |
ControlError |
Status |
|---|---|---|
Config(e), e.kind() is InvalidData or InvalidInput: a config or subscribe file that does not parse, a subscribe file that does not merge, a value lowering refuses, an [api] error |
Invalid(e.to_string()) |
422 |
Config(e), any other kind: a file that cannot be read |
Failed(e.to_string()) |
500 |
Apply(ApplyError::Bind { .. }), Apply(ApplyError::Stopped) |
Failed |
500 |
Apply(e), every other ApplyError: a duplicate tag, an unknown reference, an empty or unprobeable balancer, an invalid resource, a failed build, a disruptive change |
Invalid |
422 |
The texts are the ones etemenanki-app --test prints after configuration invalid: , and the supervisor’s ApplyError texts (Validation and apply errors).
A read failure reaches the client with the prefix Instance::reload gives it, cannot read <config path>: <error>, which keeps the error’s kind. GET and PUT read through ConfigFile::read instead, which prefixes the bare path, so the same missing piece reads differently from the two sides:
| Failure | GET, PUT before the edit |
POST /v1/reload, the reload after a PUT |
|---|---|---|
[subscribe] without path |
422, <path>: [subscribe] needs a path naming the subscribe file |
422, cannot read <path>: [subscribe] needs a path naming the subscribe file |
| Subscribe file missing | 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) |
Every reload the API asks for is logged as any reload is, under the target etemenanki_app::instance: config reloaded: <summary> at info, or at error one of reload: <error> (the files cannot be read), reload: <error>; keeping the running config, and reload refused, keeping the running config: inbound <tag>: <reason>, which would end its live connections; restart to apply it. The full list is on etemenanki-app: running, reloading and shutting down.
Connections and traffic
Section titled “Connections and traffic”The five tracking handlers in webclient/src/tracking.rs take the Tracker as their only state. router() gets it from the app as instance.tracker(), a clone of the supervisor’s, so they see every flow of every inbound and outbound, across reloads. How flows are registered, killed, announced and sampled is on Tracking flows, stats and speed limits; this section covers the HTTP shapes.
Times are Unix milliseconds, computed by unix_ms(started) as SystemTime::now() - started.elapsed() at response time: the tracker keeps a monotonic Instant, so a wall-clock step moves the reported start. A time before the epoch is 0; one past u64::MAX milliseconds saturates.
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 reads tracker.flows() first and tracker.sessions() after it: two separate reads of two structures, not one snapshot. A flow’s session can be missing from sessions (the session ended between the reads, or opened after them), and a listed session can have no flow.
sessions is Tracker::sessions(), sorted by session id. Each is a SessionView:
| Field | From SessionStats |
Meaning |
|---|---|---|
id |
id.get() |
The session id |
inbound |
inbound |
The tag of the inbound that accepted it |
source |
source |
The client’s IP, or null |
user |
user_label when user is set |
The label the user is known by, or null for nobody. A session names its user once its first flow bound it. |
up, down |
up, down |
Wire bytes from and to the client, protocol overhead included (Per-user usage accounting) |
started_at |
started |
When it was accepted |
flows is Tracker::flows(), sorted by flow id with sort_unstable_by_key. Each is a FlowView, built from a FlowEntry:
| Field | Meaning |
|---|---|
id |
The flow id: unique for the supervisor’s life, and growing in registration order |
session |
The session carrying it, or null for a flow no session carries (a TUN device’s) |
network |
"tcp" or "udp" from the destination’s DialNetwork; "unknown" for Unknown and Unix |
inbound |
The inbound tag |
user |
The user’s label, or null for nobody |
source |
The source IP, or null |
host, port |
The destination: a domain as requested, or an IP. For a UDP sub-link, the destination of the packet that opened it |
sniffed |
The domain sniffing recovered, or null; always null for a UDP sub-link |
outbound, outbound_version |
The tag and version of the outbound it was opened on. A balancer’s flow names the member. The supervisor’s own DNS service has the empty tag "". |
rule |
The index of the route rule that matched, or null for the default route and for a query the DNS service answered |
plane_epoch |
The epoch of the routing plane that routed it (The plane: routing each flow) |
started_at |
When it registered, which is after its dial completed |
up, down |
Payload bytes toward the destination and back |
The example is shaped like the output of connections_lists_live_sessions_and_flows. Its ids, port and times are made up (the test’s echo server binds port 0); the byte counts are the ones the test pins: one SOCKS connection echoing 5 bytes counts 3 + 10 + 5 = 18 wire bytes up (greeting, request, payload) and 2 + 10 + 5 = 17 down on the session, and 5 payload bytes each way on the flow.
DELETE /v1/connections/{flow}
Section titled “DELETE /v1/connections/{flow}”close_one calls Tracker::kill(FlowId::new(flow)):
true:204with an empty body.killreturnstruewhenever the id is still registered, including a flow already killed whose connection has not dropped it yet.false:404,{"error": "no live flow <flow>"}.
A kill closes that flow alone. For a mux sub-flow, the carrier and its other sub-flows keep running. For a UDP sub-link, the association and its other sub-links do, and a later packet routed to that outbound opens a new sub-link: a new flow with a new id. The table of what each protocol does with a killed flow is on Tracking flows, stats and speed limits. deleting_a_flow_closes_it_alone pins that the other connection keeps echoing and that a second DELETE after the close is 404.
DELETE /v1/connections with a selector
Section titled “DELETE /v1/connections with a selector”#[derive(Deserialize)]#[serde(deny_unknown_fields)]pub(crate) struct Selector { outbound: Option<String>, inbound: Option<String>, session: Option<u64>, user: Option<String>,}close_matching calls Tracker::kill_where(|f| selector.matches(f)) and answers 200 with {"closed": <n>}, the number of flows this call killed. A flow matches when it matches every field given:
| Field | Matches a flow whose |
|---|---|
outbound |
outbound tag equals it, whatever the version. A balancer’s tag matches nothing, since flows name the member. |
inbound |
inbound tag equals it |
session |
session id equals it. Flows without a session never match. |
user |
user is set and its label equals it. Anonymous flows never match. |
With no field, every live flow matches ("no selector closes everything"). An empty value is still a value: ?outbound= matches only the empty tag. A field that is not one of the four is 400, since Selector is deny_unknown_fields, and so is a repeated field or a session that is not a u64. kill_where skips flows that are already killed, so the same selector sent twice closes nothing the second time.
GET /v1/connections/events
Section titled “GET /v1/connections/events”events wraps tracker.subscribe() in a tokio_stream::wrappers::BroadcastStream and maps each item to an SSE event:
| Item | Event | Data |
|---|---|---|
Ok(FlowEvent::Opened(flow)) |
opened |
The FlowView, with its counts read when the stream sends the event (usually 0 for a fresh flow) |
Ok(FlowEvent::Closed(flow)) |
closed |
The FlowView with its final counts |
Err(BroadcastStreamRecvError::Lagged(missed)) |
lagged |
{"missed": <n>} |
The events carry the shared Arc<FlowEntry>, not a copy: Tracker::open sends Opened with the entry it registers, and the flow handle’s Drop sends Closed after the flow has left the registry. events builds each FlowView in the stream’s map, when the event is serialised for this client, so up, down and started_at are read then. A client that reads an opened event late sees what the flow has moved since, and for a flow that has closed meanwhile its final counts. closed always carries the final counts, since the handle that counted is gone by then.
The stream starts from the subscription: a client can see closed for a flow whose opened it never saw. The tracker’s broadcast channel keeps the last EVENT_CAPACITY = 1024 events; a client that falls further behind gets one lagged event with the number it missed and continues from the oldest event still kept. Opening and closing flows never wait for an SSE client.
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 wraps tracker.stats(), the sampler’s watch receiver, in WatchStream::from_changes, which skips the snapshot current at subscription and yields each one published after it. Each becomes a traffic event:
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 } }}| Field | From StatsSnapshot |
Meaning |
|---|---|---|
tick |
tick |
Ticks since the supervisor started |
interval_ms |
interval in milliseconds, saturating |
The time the rates are over, as measured since the previous tick |
up, down, up_rate, down_rate, flows |
total, flattened into the top level |
Payload bytes since start, bytes per second over the tick, flows live at the tick |
inbounds, outbounds |
By tag, sorted | The same five figures per inbound and per outbound tag that has carried a flow since start |
StatsSnapshot::users is left out: per-user figures are the Rust API’s alone. The etemenanki-app supervisor samples every DEFAULT_SAMPLE_INTERVAL = 1 s, so a client gets one event a second. A watch channel keeps only the latest snapshot, so a client that reads slower than the tick sees tick skip rather than a backlog.
Event-stream framing
Section titled “Event-stream framing”Both streams are axum::response::sse::Sse responses with Content-Type: text/event-stream and Cache-Control: no-cache. Each event is event: <name> and one data: line of JSON, built by json_event, which expects serialisation to succeed (“these views always serialise”). KeepAlive::default() sends a : comment line whenever 15 seconds pass without an event.
Neither stream ends by itself while the API runs:
/v1/connections/eventsnever ends while its connection lives. The broadcast sender sits in the tracker’s shared state, and the router’s own state holds a clone of the tracker./v1/trafficends when the sampler’swatchsender is dropped, which happens when the sampler stops at supervisor shutdown.
Both are dropped with their connection.
The API in etemenanki-app
Section titled “The API in etemenanki-app”From [api] to ApiOptions
Section titled “From [api] to ApiOptions”app/src/config.rs declares the section:
#[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 is Option<ApiConfig>: without an [api] section the API does not listen, and a config with only [api] listens on 127.0.0.1:9090 without a secret. api::options(cfg) parses listen with Listen::from_str, wraps the secret in Secret::new, calls ApiOptions::new, and turns an OptionsError into an io::Error of kind InvalidInput prefixed [api] . instance::build calls it next to lower, so --test, a start and every reload refuse a bad [api] like any other config error, and a reload that refuses it leaves the running API alone. The table gives the texts as etemenanki-app --test prints them. Configuration OK. is a plain line on standard output; configuration invalid: <error> is an ERROR tracing line (timestamp, level, target etemenanki_app), and the table quotes its message:
[api] |
Printed |
|---|---|
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 = "", any 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 |
An unknown key, such as cors = true |
configuration invalid: TOML parse error at line … |
[api] alone, listen = "[::1]:9090", listen = "./api.sock", or listen = "0.0.0.0:9090" with a secret |
Configuration OK. |
--test checks [api] but binds nothing, so it does not find a port in use.
Starting
Section titled “Starting”Core keeps the running config’s options in a tokio::sync::watch::Sender<Option<ApiOptions>>, created by Core::start with the first config’s options. After a successful apply, Core::reload_with updates it with send_if_modified, which notifies receivers only when the new options differ from the old. A reload that is refused, or finds nothing changed, leaves it alone. Instance::api() hands out a receiver (etemenanki-app: running, reloading and shutting down).
main starts the API right after Instance::start:
options = instance.api();first = options.borrow().clone().- With
Some(first),start_api(instance.clone(), &first). A failure is logged asfailed to start: API on <listen>: <error>, the instance is shut down withDuration::ZERO, and the process exits with failure: an[api]that cannot bind at start fails the start, as an inbound would. - It creates the
stop_apioneshot and spawnsfollow_api(instance, options, running, api_stopped). - It starts the file watcher.
start_api binds ApiListener, logs API listening on <addr> (the TCP address from local_addr; when that is None, the configured listen, which is the unix path, or a TCP address whose socket could not report its own), builds router(ConfigFile::new(instance), instance.tracker(), options.secret.clone(), env!("CARGO_PKG_VERSION")), creates a oneshot whose receiver is the graceful-shutdown future, and spawns listener.serve(app, …). It returns Api { stop, task }.
A config without [api] binds nothing beyond its inbounds (a_config_without_api_binds_nothing_new).
Following reloads: follow_api
Section titled “Following reloads: follow_api”stateDiagram-v2 [*] --> Off: start without an api section [*] --> On: start, the API binds [*] --> Failed: start, the API cannot bind Failed --> [*]: exit with failure On --> Stopping: api section changed or removed Stopping --> On: new options bind Stopping --> Off: section removed, or the new bind fails Off --> On: api section added or changed, and it binds Off --> Off: new options fail to bind On --> [*]: shutdown, after Api stop Off --> [*]: shutdown
follow_api loops on a tokio::select! of the stopped oneshot and options.changed():
stoppedfires, or thewatchsender is gone: leave the loop, stop the running API, return.- The options changed: take them with
borrow_and_update, stop the running API if there is one (Api::stop, below), then:- with
Some(wanted),start_api. A failure is logged asAPI on <listen>: <error>; it stays off until [api] changes. - with
None, logAPI stopped: the config no longer has [api].
- with
The select! is unbiased: when the stop and an options change are ready together, either may be taken first. Either way the loop ends with the running API stopped.
The old listener is closed before the new one binds, so the same address can be taken again: a secret change rebinds the same port (the_api_follows_the_api_section_across_reloads). The restart runs in follow_api, not in the reload that published it, so a reload made through the API itself, a POST /v1/reload after [api] was edited by hand for example, finishes its request before the API stops.
Stopping
Section titled “Stopping”const SHUTDOWN_GRACE: Duration = Duration::from_secs(5);
struct Api { stop: oneshot::Sender<()>, task: JoinHandle<io::Result<()>>,}Api::stop sends on stop, which completes the graceful-shutdown future, then waits for the serve task under tokio::time::timeout(SHUTDOWN_GRACE, …):
| Outcome | Action |
|---|---|
The task returned Ok(()) |
None |
The task returned an io::Error |
error: API: <error>. axum 0.8.9’s serve future always returns Ok(()) (see ApiListener), so this branch exists for the type and is not taken at this revision |
| The task panicked or was cancelled | error: API task failed: <JoinError> |
| 5 s passed | task.abort(), then warn: API: requests still in flight when it stopped were dropped |
At shutdown (SIGINT or SIGTERM), main logs shutting down, sends on stop_api, awaits follow_api (which stops the running API with the grace above), and only then calls instance.shutdown(SHUTDOWN_GRACE). The API stops before the proxy’s connections are given their grace.
Log lines
Section titled “Log lines”All are from the binary, target etemenanki_app:
| Level | Text | When |
|---|---|---|
info |
API listening on <addr or path> |
start_api bound the listener |
error |
failed to start: API on <listen>: <error> |
The first [api] cannot bind; the process exits |
error |
API on <listen>: <error>; it stays off until [api] changes |
A reload’s new [api] cannot bind |
info |
API stopped: the config no longer has [api] |
A reload removed [api] |
error |
API: <error> |
The serve task returned an error; not reachable with axum 0.8.9 |
error |
API task failed: <error> |
The serve task panicked |
warn |
API: requests still in flight when it stopped were dropped |
The serve task did not finish within SHUTDOWN_GRACE |
info |
shutting down |
A shutdown signal arrived |
The reloads the API asks for log under etemenanki_app::instance, as every reload does (Reloading on request). The webclient crate itself logs nothing, and axum, built without its tracing feature, logs no rejection.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
| The config file is the only way a switch reaches what runs: edit, then reload | ConfigControl::set_route never touches the supervisor; router::set calls reload |
a_switch_edits_the_file_reloads_and_returns_the_new_etag, a_switch_sends_new_flows_of_the_group_through_the_new_target |
| A target the view does not list is refused, and the file is untouched | Snapshot::pick |
an_invalid_target_is_422_and_leaves_the_file_untouched |
| An edit that would not start is refused, and the file is untouched | instance::check_bytes before write_atomically |
an_edit_that_would_not_start_is_422_and_leaves_the_file_untouched |
An unknown group is 404, and the file is untouched |
Snapshot::pick |
an_unknown_group_is_404 |
An edit against a stale If-Match is 409, carries the current ETag, and leaves the file untouched; the current tag goes through |
set_route’s tag comparison; IntoResponse adds the header |
an_edit_against_a_stale_etag_is_409 |
| A file changed while an edit was being checked is not written over | The re-read after check_bytes |
None |
The ETag header is the SHA-256 of the file’s bytes, before and after a switch |
ETag::of |
lists_each_groups_pick_and_inherit_and_the_targets, a_switch_edits_the_file_reloads_and_returns_the_new_etag |
| Only the pick changes; comments, ordering and formatting survive | toml_edit, set_string keeping decor |
comments_and_layout_survive_a_switch |
Without [subscribe], a pick is written as [route] default and the rules are kept |
edit choosing [route]; toml_edit printing the table’s header once it holds a key |
without_a_subscribe_file_the_default_route_is_the_one_pick |
| The file is never seen half-written | Temporary file, sync_all, rename, directory sync_all |
None |
| An existing file at a temporary name is never opened | create_new (O_EXCL) |
None |
| API edits apply one at a time, each reporting its own reload | Shared::edits held across set_route and reload |
None |
A reload of the files the API just applied is unchanged |
Core::reload_with comparing Sources |
a_reload_after_the_apis_own_finds_nothing_changed |
| With a secret, nothing is answered without the token, unknown paths included | The bearer layer over the merged router | a_secret_is_required_as_a_bearer_token |
| An API other hosts can reach has a non-empty secret | ApiOptions::new |
an_api_open_to_other_hosts_needs_a_secret |
[api] refuses unknown keys and host names; listen defaults to loopback; a path is a 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 |
A config without [api] binds nothing new |
main starts the API only for Some options |
a_config_without_api_binds_nothing_new |
The API follows [api] across reloads, and a secret change rebinds the same port |
send_if_modified, follow_api stopping before it starts |
the_api_follows_the_api_section_across_reloads |
| A unix socket file is removed only while it is still the one this process bound | SocketFile’s device and inode check |
None |
| The tracking endpoints need no config | They take the Tracker alone |
Every test in webclient/tests/tracking.rs runs with NoConfig, whose methods all fail |
| An SSE client never slows the data plane | Bounded broadcast and a watch channel on the tracker side |
None here; see Tracking flows, stats and speed limits |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Failure | Result | File |
|---|---|---|
| Config or subscribe file cannot be read | GET and PUT: 500 with <path>: <error>. POST /v1/reload and the reload after a PUT: 500 with cannot read <path>: <error> |
Untouched |
Config or subscribe file does not parse; [subscribe] without path |
GET and PUT: 422. The reloads: 422, the read error prefixed cannot read <path>: |
Untouched |
{group} not UTF-8 after percent-decoding |
400, plain text from axum, before the handler runs |
Untouched |
Bad body or If-Match |
400, before the edits lock is taken |
Untouched |
| Unknown group; unknown target | 404; 422 |
Untouched |
Stale If-Match; file changed during the check |
409 with the current ETag |
Untouched |
| Edited config fails the check | 422 (or 500 for an error of another kind, such as a file the lowering cannot read) |
Untouched |
| The re-read after the check fails | 500 with the bare I/O error, no path prefix |
Untouched |
The config path cannot be resolved (canonicalize), its metadata cannot be read, or no temporary file can be created |
500, writing <path>: <error>; nothing was created |
Untouched |
| The temporary file’s permissions, write, sync or rename fail | 500, writing <path>: <error>; the temporary file is removed |
Untouched |
| Directory sync after the rename fails | 500, writing <path>: <error> |
Replaced |
| Reload after a successful write fails | The reload’s error status, with the new ETag |
Replaced; the running config is the old one |
First [api] cannot bind |
failed to start: API on …, exit with failure |
— |
A reload’s new [api] cannot bind |
Logged; the API is off until [api] changes |
— |
Cancellation:
-
A
PUThandler dropped mid-request. The handler is a future run by its connection’s task. It is dropped when the client closes the connection before the answer (hyper sees the end of the stream on the busy connection and ends it), or when the runtime shuts down aftermainreturns. Dropping it releasesedits. What has happened by then depends on where it was waiting:Awaiting State of the file What runs edits.lock()Untouched; nothing was read Unchanged instance::check_bytesinsideset_routeUntouched; the check’s work is dropped Unchanged control.reload()Written: write_atomicallyis synchronous, with no await, so once it starts it runs to the endThe reload stops where it was awaited; an apply the supervisor’s actor has already started is not stopped by that (etemenanki-app: running, reloading and shutting down) -
An SSE client that disconnects. The stream is dropped, and with it the broadcast receiver or the
watchreceiver; nothing on the tracker side waits for it. -
Stopping the API. The graceful shutdown stops accepting and releases the listening socket at once, then waits for connections to finish;
Api::stopbounds that wait atSHUTDOWN_GRACEand then aborts the serve task. -
Stopping the process. The API stops first, then
Instance::shutdown(SHUTDOWN_GRACE)gives connections their grace (etemenanki-app: running, reloading and shutting down).
Limits
Section titled “Limits”| Item | Value | Defined in | Meaning |
|---|---|---|---|
DEFAULT_LISTEN |
127.0.0.1:9090 |
webclient/src/listen.rs |
listen when [api] does not set it |
DEFAULT_GROUP |
"default" |
webclient/src/lib.rs |
The group name that stands for the default route |
SHUTDOWN_GRACE |
5 s | app/src/main.rs |
How long Api::stop waits for the serve task, and the grace Instance::shutdown gets at exit |
| Temporary names | 64 attempts | create_beside |
Names tried before no free temporary name for … |
| Temporary file mode | 0o600, narrowed by the umask |
create_beside |
While the file is empty, until the original’s mode is copied over |
| Request body | 2,097,152 bytes | axum’s default body limit | Larger PUT bodies are 400 |
| HTTP versions | HTTP/1 | Workspace axum features (http1 only) |
No HTTP/2 |
| SSE keep-alive | A : comment after 15 s without an event |
axum’s KeepAlive::default() |
What an idle stream carries |
EVENT_CAPACITY |
1024 | supervisor/src/track/mod.rs |
Events a /v1/connections/events client may fall behind before lagged |
| Sample interval | 1 s | DEFAULT_SAMPLE_INTERVAL, kept by etemenanki-app |
One /v1/traffic event per second |
| ETag | 64 lower-case hex digits, quoted | etag.rs |
SHA-256 of the config file only |
Unit tests of the app’s side in app/tests/unit/api.rs (the tests module of app/src/api.rs). The async tests each start a real Instance on a config in its own temporary directory (etemenanki-api-<pid>-<n>) with no inbounds, build the router over ConfigFile and instance.tracker() with the version "test", and drive it with oneshot. The two [api] tests are synchronous: they call api::options on a parsed section (api_section), with no instance and no router. The fixture is a client config with comments everywhere an edit could lose them, and a subscribe file with the nodes us-socks and jp-socks and the groups Video Streaming (inheriting Proxy) and Proxy.
| Test | Behaviour it pins |
|---|---|
lists_each_groups_pick_and_inherit_and_the_targets |
The whole RouteView JSON: picks, inherit as the subscribe file writes it, targets resolved through the chain, target order and kinds; the ETag header is the file’s |
a_switch_edits_the_file_reloads_and_returns_the_new_etag |
200, "status": "applied", the new ETag equals the file’s; the file has the pick; a following GET shows it with the same tag |
the_default_route_is_switched_through_its_own_name |
PUT /v1/routes/default writes [subscribe.routes] default |
an_invalid_target_is_422_and_leaves_the_file_untouched |
422 naming the target; the bytes are unchanged |
an_edit_that_would_not_start_is_422_and_leaves_the_file_untouched |
A known target (direct) that is a balancer fails the check: 422 containing names a balancer; the bytes are unchanged |
an_unknown_group_is_404 |
404; the bytes are unchanged |
an_edit_against_a_stale_etag_is_409 |
After a hand edit, the old tag is 409 with the current tag and the hand edit intact; the current tag goes through |
a_secret_is_required_as_a_bearer_token |
No header is 401 with WWW-Authenticate: Bearer; a token one character short is 401; an unknown path without a token is 401; the right token is 200 |
comments_and_layout_survive_a_switch |
Replacing a pick keeps its line comment; a new pick joins its table; the rest of the file is byte for byte the same |
a_reload_after_the_apis_own_finds_nothing_changed |
After a switch, POST /v1/reload is {"status": "unchanged"} |
without_a_subscribe_file_the_default_route_is_the_one_pick |
No groups; default target is the first outbound; a switch writes [route] default and keeps the rules; picking direct adds no outbound; any other group is 404 |
api_options_come_from_the_api_section |
No [api] is None; [api] alone is DEFAULT_LISTEN without a secret; absolute and relative paths are unix sockets; a non-local listen with a secret is accepted |
an_api_open_to_other_hosts_needs_a_secret |
0.0.0.0 without a secret fails with needs a secret; an empty secret fails; a host name fails; an unknown key fails |
End-to-end tests in app/tests/integration/e2e_api.rs run the app binary and speak raw HTTP/1.1 to it:
| Test | Behaviour it pins |
|---|---|
a_config_without_api_binds_nothing_new (Linux only) |
From /proc/<pid>/fd and /proc/<pid>/net/{tcp,tcp6,unix}: without [api] the process listens on its SOCKS port alone and no unix socket; with [api] on exactly that port and the API’s |
a_switch_sends_new_flows_of_the_group_through_the_new_target |
A group inheriting blackhole drops loopback traffic; after PUT to a node (a second app), new connections reach the echo server, and stop reaching it when the node stops; after PUT to direct they reach it again; GET and the file show the pick |
the_api_follows_the_api_section_across_reloads |
/v1/health answers; removing [api] stops it; adding it back on another port with a secret starts it (401 without the token); changing only the secret rebinds the same port, and the old secret is 401 |
Tracking tests in webclient/tests/tracking.rs run a real supervisor with a SOCKS inbound sending everything to a freedom outbound direct, and the router over NoConfig:
| Test | Behaviour it pins |
|---|---|
connections_lists_live_sessions_and_flows |
One session and one flow with the fields above; session wire bytes 18 and 17, flow payload 5 and 5; rule is null; started_at is set |
deleting_a_flow_closes_it_alone |
204; that connection closes; the other keeps echoing; a second DELETE is 404 |
deleting_by_selector_closes_the_matching_flows |
An unmatched outbound closes 0; an unknown field is 400; outbound plus inbound closes both flows; no selector closes everything |
traffic_streams_one_event_per_tick |
With a 100 ms interval, text/event-stream; consecutive tick values; five ticks take between 4.5 and 10 intervals; the 1000 bytes sent show up in up |
connection_events_stream_opens_and_closes |
opened naming the outbound, then closed with the same id and final counts of 7 and 7 |
ffi/tests/proxy.rs → set_route_and_reload_switch_like_the_rest_api exercises routes::Snapshot through the mobile library (Mobile library (etemenanki-ffi)).
No test covers: If-Match parsing (*, weak tags, lists, 400); the re-read after the check; write_atomically and create_beside (symlinks, permissions, temporary names); the lagged event; ApiListener on a unix socket and SocketFile; the 400 answers for a malformed body or a {group} that is not UTF-8; a PUT whose pick is already the target, or is written as a literal string; /v1/version; a reload that fails after a write; and Api::stop’s timeout. Add one when you touch those paths.
Run them with:
cargo test -p etemenanki-app --lib api::testscargo test -p etemenanki-app --test integration e2e_apicargo test -p etemenanki-webclientFor the test layout and harnesses, see Testing.