katana internals
Source files: 43 · checked against katana v3.0.1 · Etemenanki 596916d
katana/Cargo.tomlkatana/src/main.rskatana/src/runtime.rskatana/src/config.rskatana/src/api/mod.rskatana/src/api/newv2board.rskatana/src/api/sspanel.rskatana/src/manager/mod.rskatana/src/manager/node.rskatana/src/manager/transport.rskatana/src/manager/proxy.rskatana/src/serve.rskatana/src/connector.rskatana/src/meter.rskatana/src/traffic.rskatana/src/router.rskatana/src/rule.rskatana/src/inbound.rskatana/src/outbound/mod.rskatana/src/outbound/freedom.rskatana/src/outbound/proxy.rskatana/tests/integration.rskatana/tests/support/mod.rskatana/tests/unit/e2e.rskatana/tests/unit/serve.rskatana/tests/unit/connector.rskatana/tests/unit/traffic.rskatana/tests/unit/meter.rskatana/tests/unit/inbound.rskatana/tests/unit/runtime.rsEtemenanki/concepts/src/link.rsEtemenanki/concepts/src/runtime.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/flow.rsEtemenanki/protocols/src/transports/accept.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/wireguard/device.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/environment/src/routing.rsEtemenanki/app/src/serve.rsEtemenanki/app/src/connector.rsEtemenanki/app/src/flow.rsEtemenanki/app/src/instance.rs
katana is a panel node agent. One process serves one or more panel nodes, keeps each node’s user list in step with its panel, and bills and paces every user’s traffic. It lives in its own repository and is a separate program from etemenanki-app. The protocol machinery comes from the Etemenanki kernel. katana builds everything that makes it a node agent: the panel clients, the node lifecycle, its own serve loop, flow admission, metering, traffic accounting and destination audit.
This page is the map for contributors. It says which parts come from the kernel and which live in katana, lists every module, draws the task tree, names the main data structures and their owners, and follows one request through katana next to the same request in etemenanki-app. Each section links to the page that covers that part in detail.
Responsibilities
Section titled “Responsibilities”katana depends on three kernel crates, taken from a private Cargo registry. Cargo.toml requires 2.0.0 for each, and Cargo.lock resolves etemenanki-concepts and etemenanki-environment to 2.0.0 and etemenanki-protocols to 2.0.1. Through etemenanki-protocols 2.0.1 katana gets the bounded per-connection uplink in the WireGuard driver, which makes a slow tunnel push back on the writer, and the mux fixes that send each downlink frame once and keep every frame of a VMess read. It builds etemenanki-protocols with the features vendored-openssl and hysteria. katana does not depend on etemenanki-app. Its configuration, serve loop, connector and outbound pool are its own code, written in the same shape as the app’s.
Taken from the kernel
Section titled “Taken from the kernel”| Kernel crate | What katana uses | Used in |
|---|---|---|
etemenanki-concepts |
The server runtime ProxyServerRuntime in progress mode (showing_progress) and its Traffic deltas; ProxyCoreDecode; the Connector and DatagramLink traits and link::Outbound; Destination, Remote, DialNetwork, NetworkUser; SniffedBehavior; the client runtime ProxyClientConnector and ProxyClientRuntime with ProxyCoreEncode, ProxyCoreEncodeDatagram and NoCodec |
src/serve.rs, src/connector.rs, src/router.rs, src/outbound/ |
etemenanki-environment |
The route model routing::{Router, RouteTable, RouteItem, RouteMatch, RouteTarget, build_geo_data, parse_port_match}; the dialers dial::{Dialer, SocketOptions, DualStackUdp} |
src/router.rs, src/outbound/ |
etemenanki-protocols |
The server cores VMessCore, VlessCore, TrojanCore, ShadowsocksCore, Ss2022Core and their user validators; Flow; HANDSHAKE_TIMEOUT and RELAY_IDLE_TIMEOUT; InboundTransport, TransportStream, TransportConnector and the TLS ServerConfig; the Hysteria 2 listener Hy2Inbound with Authenticator and Masquerade; the client codecs VMessStream, VMessDatagram, VlessStream, VlessDatagram, SsStream, Ss2022Stream, HttpConnect, SocksConnect and SocksUdpLink; the WireGuard WgConnector; the DNS Resolver; AddressFamilyStrategy |
src/serve.rs, src/connector.rs, src/inbound.rs, src/manager/, src/outbound/, src/runtime.rs |
The kernel owns every byte on the wire: framing, authentication inside the protocol handshake, sniffing, transports, the per-connection runtime, and the router’s matching engine. katana never parses a proxy protocol itself. Every SOCKS, HTTP, VMess, VLESS and Shadowsocks outbound is dialed through a TransportConnector of TransportKind::Tcp, so katana builds no TLS, WebSocket or gRPC transport toward an upstream. A WireGuard outbound runs its own tunnel through WgConnector and uses no TransportConnector.
Implemented in katana
Section titled “Implemented in katana”| Concern | Module | Nearest equivalent in etemenanki-app |
|---|---|---|
Panel clients for UniProxy (Xboard, V2board) and SSPanel mod_mu, with ETag caching |
src/api/ |
None; the app has no panel |
| Process root, config watcher, node add, remove and reconfigure | src/runtime.rs |
app/src/instance.rs |
| Node lifecycle: bootstrap with retry, poll cycle, reconcile ladder, static updates | src/manager/node.rs |
None; the app replaces whole generations |
| One listener generation and its task scope | src/manager/transport.rs, src/serve.rs |
app/src/serve.rs |
| Swappable user tables and the pre-auth limit | src/manager/proxy.rs |
Fixed tables per generation |
| Flow admission and user retirement | src/connector.rs → Admission |
None |
| Routing and audit per flow and per UDP packet | src/connector.rs, src/router.rs, src/rule.rs |
app/src/connector.rs → AppConnector (routing only) |
| Metering and the per-user speed limit | src/meter.rs, src/traffic.rs → TokenBucket |
None |
| Traffic accounting with draining counters and residuals | src/traffic.rs → NodeTraffic |
None |
| Outbound pool | src/outbound/ |
app/src/outbound/ |
Module map
Section titled “Module map”Sizes are line counts at v3.0.1, comments included. The source files add up to about 6,460 lines and the tests to about 4,400.
| File | Lines | Responsibility |
|---|---|---|
src/main.rs |
54 | clap Args (-c/--config, default config.toml; --test), tracing setup, dispatch to runtime::test_config or runtime::run. Includes the end-to-end test module. |
src/runtime.rs |
439 | Process root: loads the config, builds the outbound pool and the shared DNS resolver, builds (build_node) and spawns (spawn_built) one node task per [[node]], runs the config watcher and apply_reload, handles signals. Also test_config and init_tracing. |
src/config.rs |
324 | The serde schema, with deny_unknown_fields on every struct; load and parse_bytes. |
src/api/mod.rs |
369 | Panel-neutral models (NodeType, Transport, NodeInfo, UserInfo, DetectRule, DetectResult), EtagCache, error_for_status without the URL, the local rule-file reader, panel_node_type, and the PanelClient enum with its constructor PanelClient::new, forget_etags and inherit_routes. |
src/api/newv2board.rs |
452 | The UniProxy client and its parsers, and node_type_param. |
src/api/sspanel.rs |
576 | The mod_mu client, the custom_config and legacy parsers, compare_version. |
src/manager/mod.rs |
97 | StaticUpdate, user_key, user_tag, node_tag, build_user_entries, user_set_differs. |
src/manager/node.rs |
675 | NodeManager: bootstrap and its retry with backoff, poll cycle, reconcile ladder, static updates, rule refresh, traffic and audit reports. |
src/manager/transport.rs |
177 | TransportManager: one bound listener generation, start, shutdown, and the Drop backstop. |
src/manager/proxy.rs |
239 | ProxyManager: Tables, the pre-auth semaphore, the handshake-failure sampler, accept_stream and refresh. |
src/serve.rs |
447 | Scope, spawn_scoped, accept_loop, serve_socket, serve_stream, drive, until_retired, run_hysteria. |
src/connector.rs |
361 | Admission, Dispatcher, KatanaConnector, LeaseSlot, and the UDP FanOut. |
src/meter.rs |
141 | Gate and the Metered<S> stream wrapper. |
src/traffic.rs |
472 | AuthKey, UserTag, TokenBucket, UserCounter, NodeTraffic, PreparedUsers, determine_rate. |
src/router.rs |
120 | route_target and build_router over the kernel route model. |
src/rule.rs |
76 | RuleManager: audit rules per node tag and the deduplicated hit ledger. |
src/inbound.rs |
548 | StreamProtocol, build_transport, build_protocol, the Shadowsocks tables, the Hysteria authenticator and listener builders, validate_hysteria. |
src/outbound/mod.rs |
533 | Outbound, connect_stream, connect_datagram, build_outbound, SocksOutbound, and the Shadowsocks, WireGuard and address-family builders. |
src/outbound/freedom.rs |
181 | FreedomConnector (direct TCP and dual-stack UDP) and ResolvingUdp. |
src/outbound/proxy.rs |
179 | ProxyClient, OutboundStream, OutboundDatagram. |
Outside src/api/, only two modules touch a panel client, and both build it with PanelClient::new from src/api/mod.rs: src/runtime.rs builds each node’s first client and checks a reloaded node’s client, and src/manager/node.rs builds a replacement in apply_static and makes every panel call. The builders in src/inbound.rs and src/manager/ read the api models, and nothing on the data path makes a panel request.
Task tree
Section titled “Task tree”Every long-lived task hangs off one of two roots: the process root token, or the Scope of a listener generation.
flowchart TB main["runtime::run, the main task"] sig["signal task"] bridge["watcher bridge, a std thread"] node["NodeManager::run, one per node"] tm["TransportManager and its accept Scope"] mon["auth monitor, every 1 s"] acc["accept_loop"] sock["serve_socket, one per socket"] strm["serve_stream and drive, one per stream"] hy["run_hysteria"] hyk["Hy2Inbound::run and its kernel tasks"] sig -- "cancels root" --> main bridge -- "unit ticks" --> main main -- "tokio::spawn with root.child_token()" --> node node -- "owns while bound" --> tm tm --> mon tm -- "stream node" --> acc acc --> sock sock -- "ProxyManager::accept_stream" --> strm tm -- "Hysteria 2 node" --> hy hy --> hyk
runtime::runruns on the main task and owns the loaded config. It spawns oneNodeManager::runper node withtokio::spawn, giving each a child of the root token and the receiving end of aStaticUpdatechannel, and keeps aNodeHandleper node.- The signal task waits in
wait_for_shutdownforSIGINTorSIGTERMand cancels the root token. - The watcher bridge is a
std::thread. It forwards everynotifyevent on the config file’s parent directory into a tokio channel as(). The main task sleeps 500 ms to debounce, drains the channel withtry_recv, reparses the file and callsapply_reload. NodeManager::runis one task per node. It first runsbootstrap, which races each attempt against the node token and, after a failure, waits in its ownselect!over the node token, a backoff timer and theStaticUpdatereceiver. Once the node is up,serveruns the steady state: aselect!over the node token, the pollintervaland theStaticUpdatereceiver. The panel poll and the static updates share this one task, so a node never reconciles two changes at once.- Everything below the node runs in the listener generation’s
Scope.spawn_scopedwraps each future in aselect!against the scope’s token and registers it with the scope’sTaskTracker. The auth monitor, the accept loop, everyserve_socketand every per-stream task are scoped. On a Hysteria 2 node the scope holdsrun_hysteria. The kernel’sHy2Inbound::rungets the scope token and runs its own runtime per proxy stream and, when[node.hysteria].udpis on, one per connection for UDP.
Channels and tokens
Section titled “Channels and tokens”| Channel or token | Type | From → to | Purpose |
|---|---|---|---|
root |
CancellationToken |
signal task → main loop and every node | Process shutdown. |
NodeHandle.shutdown |
child CancellationToken |
apply_reload → one node |
Removes a node on reload; also fires with root. |
static_tx / static_rx |
mpsc::channel(16) of StaticUpdate |
apply_reload → one node |
A changed [[node]] table or a rebuilt outbound pool. |
ev_tx / ev_rx |
mpsc::unbounded_channel::<()>() |
watcher bridge → main loop | “Something in the config directory changed”. |
Scope.token |
CancellationToken |
TransportManager → every scoped task |
Tears one listener generation down. |
| Lease | CancellationToken per AuthKey |
Admission → every Gate and connection of that user |
Retires a user. |
LeaseSlot |
watch::Sender<Option<CancellationToken>> |
KatanaConnector → the connection’s drive |
Tells a connection whose lease it carries, once its first flow is admitted. |
Shutdown order
Section titled “Shutdown order”A listener generation is torn down the same way whether the process is stopping, the node was removed from the config, or the listener is being rebuilt: tear_down takes the TransportManager out of its Option and awaits TransportManager::shutdown. The diagram shows the full node shutdown. A rebuild runs only the tear_down part and then calls bring_up; it sends no report.
sequenceDiagram participant R as runtime::run participant N as NodeManager::run participant T as TransportManager participant S as accept Scope participant A as Admission participant P as Panel R->>N: node token cancelled N->>T: tear_down takes the Option, calls shutdown T->>S: Scope::shutdown cancels, closes, waits S-->>T: every scoped task has finished T->>A: retire_all cancels every lease T->>T: release_listener awaits Hy2Inbound::shutdown N->>P: report_traffic, then report_illegal R->>N: awaits the JoinHandle
TransportManager::shutdown returns only after every task of the generation is gone, every lease is cancelled, and, on a QUIC node, the UDP port is free again. The last step exists because quinn releases its socket only after its driver has been polled with no connections left, and NodeManager::rebuild binds the same port straight after. The kernel bounds that wait: Hy2Inbound::shutdown waits up to DRAIN_TIMEOUT (3 s) for the endpoint to go idle, then tries to rebind the port every RELEASE_POLL (20 ms) for up to RELEASE_TIMEOUT (3 s), and logs a warning if it never comes free.
Once every relay has stopped, the counters are final, so the node sends one last traffic and audit report before its task returns.
impl Drop for TransportManager is a backstop for a generation dropped without shutdown: it cancels the scope token and calls retire_all. A Drop cannot wait for the drain, which is why every path calls shutdown explicitly.
Key data structures
Section titled “Key data structures”Process level
Section titled “Process level”pub type LogReload = Box<dyn Fn(&str) + Send + Sync>;
type Pool = Arc<HashMap<CompactString, Arc<Outbound>>>;
type NodeId = (String, String, u32, String, String);
struct NodeHandle { id: NodeId, cfg: NodeConfig, task: JoinHandle<()>, static_tx: mpsc::Sender<StaticUpdate>, shutdown: CancellationToken,}
pub async fn run(config_path: PathBuf, reload: LogReload) -> ExitCodefn spawn_node(cfg: &NodeConfig, pool: &Pool, root: &CancellationToken) -> Option<NodeHandle>fn build_node(cfg: &NodeConfig, pool: &Pool) -> anyhow::Result<Arc<NodeManager>>fn spawn_built(nm: Arc<NodeManager>, cfg: &NodeConfig, root: &CancellationToken) -> NodeHandlepub fn build_outbounds(cfg: &Config) -> io::Result<HashMap<CompactString, Arc<Outbound>>>Poolis the process’s one outbound pool.build_outboundsbuilds oneResolverfrom[dns]and shares it with every outbound. It seeds the reserved tagsdirectandfreedom(aFreedomConnectorwithAddressFamilyStrategy::Auto) andblockandblackhole(Outbound::Block), then adds each[[outbound]]. A duplicate or reserved tag fails withduplicate/reserved outbound tag <tag>. Each node holds the pool’sArc, and each compiledRouterholdsArc<Outbound>clones of the entries its rules name.build_nodebuilds a node’s panel client withPanelClient::newand itsNodeManager, which compiles the router; nothing is bound yet.spawn_builtcreates the node’s channel and child token and spawnsNodeManager::run.spawn_nodejoins the two at startup and logs and skips a node that does not build.NodeIdis the panel node an entry serves:(panel_type lowercased, api.host, api.node_id, api.key, panel node type). The last field ispanel_node_type: on newV2board it is thenode_typeUniProxy is asked for (node_type_param:vlessfor a V2ray-family node,node_typeV2ray,VmessorVless, withenable_vless, else the lowercasednode_type), because UniProxy finds a node by its id and that type; on SSPanel it is empty. A change to any field respawns the node with a fresh traffic registry.display_idrenders the identity for logs astype@host#node_id, with/node_typeappended on newV2board, and never showsapi.key.apply_reloadmatches nodes byNodeId. Before it touches a running node, it builds every node the reload adds (build_node, against the new pool if the outbounds changed) and builds the panel client of every kept node whoseNodeConfigchanged. If any of them fails, it logsreload: node <id>: <error>; keeping current configand applies nothing. Otherwise it swaps the pool and sendsStaticUpdate::Outbounds, cancels and awaits each node whose identity is gone, spawns the pre-built nodes withspawn_built, and sendsStaticUpdate::Configto each kept node whoseNodeConfigdiffers.
pub enum StaticUpdate { Config(Box<NodeConfig>), Outbounds(Arc<HashMap<CompactString, Arc<Outbound>>>),}Node level
Section titled “Node level”pub struct NodeManager { api: Mutex<Arc<PanelClient>>, cfg: Mutex<NodeConfig>, pool: Mutex<Arc<HashMap<CompactString, Arc<Outbound>>>>, router: Mutex<Arc<Router<Outbound>>>, traffic: Arc<NodeTraffic>, rules: Arc<RuleManager>, transport: tokio::sync::Mutex<Option<TransportManager>>, cur: tokio::sync::Mutex<Cur>,}
impl NodeManager { pub fn new( api: PanelClient, cfg: NodeConfig, pool: Arc<HashMap<CompactString, Arc<Outbound>>>, ) -> io::Result<Arc<Self>>;
pub async fn run( self: Arc<Self>, shutdown: CancellationToken, mut static_rx: mpsc::Receiver<StaticUpdate>, );}| Field | Lives for | Notes |
|---|---|---|
api: Mutex<Arc<PanelClient>> |
Until a config edit replaces it | An enum over sspanel::Client and newv2board::Client, built by PanelClient::new. Each call clones the Arc out of the lock, so no guard is held across a request. A client keeps its own copy of the settings it was built from, so a StaticUpdate::Config that changes panel_type or [node.api] builds a new client, copies the newV2board routes the old one learned (inherit_routes, which keeps the audit rules derived from them in force until the new client reads the node config) and swaps it in. Each client holds a reqwest::Client whose timeout is api.timeout seconds (0 means 5 s), and an EtagCache keyed by endpoint: "node" and "users", plus "rules" on SSPanel. Every bootstrap attempt calls forget_etags first. |
cfg, pool, router |
The node task | parking_lot::Mutex, replaced by static updates. A config edit compiles its router from the new route before it stores anything and is refused if compilation fails. A new outbound pool goes through rebuild_router, which recompiles router from cfg.route and pool and keeps the old router if compilation fails. |
traffic: Arc<NodeTraffic> |
The node task | Outlives every listener generation, so byte totals continue across rebuilds. |
rules: Arc<RuleManager> |
The node task | Audit rules stored per node tag Type_listenip_port (node_tag). |
transport |
One listener generation | None while the node has no users, and after a failed rebuild. |
cur: Cur |
The node task | The last applied NodeInfo, user list and node tag. The poll cycle falls back to it when a fetch fails or returns 304. |
Listener generation
Section titled “Listener generation”pub struct TransportManager { accept: Scope, proxy: Arc<ProxyManager>,}
impl TransportManager { pub async fn start( node: &NodeInfo, cert: &CertConfig, listen_ip: &str, enable_vless: bool, sniff: bool, users: &[UserInfo], traffic: Arc<NodeTraffic>, rules: Arc<RuleManager>, router: Arc<Router<Outbound>>, node_tag: CompactString, hysteria: &HysteriaConfig, ) -> io::Result<Self>;
pub async fn shutdown(self);}#[derive(Clone, Default)]pub struct Scope { pub token: CancellationToken, tasks: TaskTracker,}
pub fn spawn_scoped<F>(scope: &Scope, fut: F)where F: Future + Send + 'static, F::Output: Send,pub enum Tables { Stream(ArcSwap<StreamProtocol>), Hysteria { server: Hy2Inbound<UserTag>, cfg: HysteriaConfig, },}
pub struct ProxyManager { tables: Tables, sniff: bool, dispatcher: Arc<Dispatcher>, traffic: Arc<NodeTraffic>, preauth: Arc<Semaphore>, handshake_failures: AtomicU64, node_tag: CompactString,}pub struct Dispatcher { pub router: Arc<Router<Outbound>>, pub rules: Arc<RuleManager>, pub node_tag: CompactString, pub admission: Admission,}
pub struct Admission { traffic: Arc<NodeTraffic>, leases: Mutex<HashMap<AuthKey, CancellationToken>>,}A listener generation fixes three things for its whole life: the bound socket, the Router inside its Dispatcher, and the kind of Tables. Only the user table inside Tables changes in place. A route edit or an outbound pool edit therefore needs a new generation. A route edit forces one through poll_cycle(true), an outbound pool edit through apply_route_change, and both end in rebuild, which by design drops every connection on the node. A user list change needs none.
Tables::Streamholds theStreamProtocol(Vmess,Vless,Trojan,ShadowsocksLegacyorShadowsocks2022) in anArcSwap.accept_streamcallsload_full()once per stream, so a connection keeps the table it authenticated against for its whole life.Tables::Hysteriaholds the running listener. A refresh swaps only its authenticator withHy2Inbound::set_authenticator; the socket and every QUIC connection stay.handshake_failuresis drained once a second by the auth monitor, which logs a warning when more thanHANDSHAKE_FAILURE_ALERT_PER_SECfailures arrived in that second. It only detects; it never blocks.
Per connection and per flow
Section titled “Per connection and per flow”pub type LeaseSlot = watch::Sender<Option<CancellationToken>>;
#[derive(Clone)]pub struct KatanaConnector { disp: Arc<Dispatcher>, source: Option<IpAddr>, lease: Option<Arc<LeaseSlot>>,}
impl Connector<Flow<UserTag>> for KatanaConnector { type Stream = Metered<OutboundStream>; type Datagram = FanOut; type Future = ConnectFuture;
fn connect(&mut self, flow: Flow<UserTag>) -> ConnectFuture;}pub enum AuthKey { Uuid(Uuid), Name(CompactString),}
pub struct UserTag { pub key: AuthKey, pub uid: i64,}
pub struct NodeTraffic { users: RwLock<HashMap<AuthKey, Arc<UserCounter>>>, residuals: Mutex<HashMap<i64, (u64, u64)>>, draining: Mutex<Vec<Arc<UserCounter>>>,}UserTagis the payload of every protocol user table, asArc<UserTag>in the kernel’sNetworkUser::user_data. It names a user; it does not carry the user’s counter. Each flow looks the counter up withNodeTraffic::lookup(key, uid), which succeeds only while that credential is registered to that same uid. A table kept alive by a long connection therefore cannot keep billing, or admitting, a user the panel has removed.AuthKeyisUuidfor VMess and VLESS nodes andName(label)for Trojan, Shadowsocks and Hysteria 2 nodes.NodeType::keys_by_emailis the single place that decides this, so the counters and the protocol tables always agree on the key.KatanaConnector: one is built per stream connection (with aLeaseSlot). On a Hysteria 2 node the kernel builds one per proxy stream and one per connection’s UDP runtime, all without aLeaseSlot.Gate(src/meter.rs) joins the user’sArc<UserCounter>and lease.Metered<S>wraps each TCP outbound with one;FanOutholds one per UDP association.
Ownership at a glance
Section titled “Ownership at a glance”flowchart TB NM["NodeManager"] PC["PanelClient"] NT["NodeTraffic"] RM["RuleManager"] TM["TransportManager"] SC["Scope"] PM["ProxyManager"] TBL["Tables"] DI["Dispatcher"] AD["Admission"] RT["Router"] UC["UserCounter"] NM --> PC NM --> NT NM --> RM NM --> TM TM --> SC TM --> PM PM --> TBL PM --> DI DI --> RT DI --> RM DI --> AD AD --> NT NT --> UC
Arrows point from owner to owned; a shared Arc is drawn from each holder. NodeTraffic and RuleManager belong to the node and are shared into each generation. Each generation’s Dispatcher gets a clone of the node’s Router Arc. Every Gate holds one more Arc<UserCounter>, and that reference count is what tells NodeTraffic::snapshot whether a draining counter still has writers.
| Lock | Type | Guards | Taken by |
|---|---|---|---|
NodeManager.api, .cfg, .pool, .router |
parking_lot::Mutex |
Panel client, static config, pool, compiled router | The node task only. |
NodeManager.transport, .cur |
tokio::sync::Mutex |
The generation, the last applied state | The node task only. |
Tables::Stream |
ArcSwap<StreamProtocol> |
The stream user table | accept_stream reads, refresh stores. |
Admission.leases |
parking_lot::Mutex<HashMap<AuthKey, CancellationToken>> |
Leases | admit, commit, retire_all. |
NodeTraffic.users |
parking_lot::RwLock |
The registry | prepare, lookup and snapshot read; commit writes. |
NodeTraffic.draining, .residuals |
parking_lot::Mutex |
Outgoing counters, counterless bytes | commit, snapshot, the report paths. |
TokenBucket.state |
parking_lot::Mutex<BucketState> |
Balance and last refill | charge, ready_at. |
RuleManager.inbound, .results |
parking_lot::RwLock, parking_lot::Mutex |
Rules per tag, hits per tag | update, detect, drain. |
ProxyClient.inner |
parking_lot::Mutex |
The kernel client connector | Held only while building a dial future. |
A parking_lot guard is not Send, and NodeManager::run is spawned with tokio::spawn, which requires a Send future. Holding one of the node’s parking_lot guards across an .await therefore does not compile. The node keeps each such lock inside one expression or a short block.
Where locks nest, they are always taken in the order Admission.leases → NodeTraffic.users → NodeTraffic.draining. admit takes leases, then users. Admission::commit takes leases, and inside it NodeTraffic::commit takes users, then draining. snapshot takes users, then draining, and takes residuals only after releasing both. Separately, RuleManager::detect takes results while it holds inbound.
Request flow
Section titled “Request flow”A TCP flow on a stream node
Section titled “A TCP flow on a stream node”sequenceDiagram participant C as Client participant L as accept_loop participant K as serve_socket participant PM as ProxyManager participant D as drive and runtime participant KC as KatanaConnector participant AD as Admission participant O as Outbound C->>L: TCP connect L->>L: try a session permit, then wait for a stage permit L->>K: spawn_scoped K->>K: InboundTransport::accept runs TLS, WebSocket or HTTP/2 K->>PM: accept_stream for each stream yielded PM->>PM: try a pre-auth permit, load_full the user table PM->>D: spawn_scoped serve_stream D->>D: the core authenticates and parses the request D->>KC: connect(flow) KC->>AD: admit(tag) AD-->>KC: counter and lease KC->>KC: publish the lease, route, audit KC->>O: connect_stream(destination) O-->>D: Metered outbound stream D->>D: established, drop the pre-auth permit D->>C: reply, then relay until end, retirement or watchdog
- Accept.
accept_looptakes a place underMAX_LIVE_CONNECTIONS_PER_NODEwithtry_acquire_owned. A socket over that limit is dropped and counted as a handshake failure. The socket then waits for a transport-stage place underMAX_TRANSPORT_STAGES_PER_NODE, andserve_socketis spawned. - Transport.
serve_socketrunsInboundTransport::accept. The stage place is released when the first stream arrives or afterTRANSPORT_HANDSHAKE_TIMEOUT, whichever comes first. A gRPC socket yields one stream per HTTP/2 stream, and all of them share the socket’s live-connection place through anArc<OwnedSemaphorePermit>. - Pre-auth.
ProxyManager::accept_streamtakes a place underMAX_PREAUTH_STREAMS_PER_NODEwithtry_acquire_owned, loads the current user table and spawnsserve_stream. A stream over the limit is dropped and counted as a handshake failure. - Handshake.
serve_streamcreates the connection’sLeaseSlotandKatanaConnector, builds the protocol core over the table, and callsdrive.drivewraps the core inProxyServerRuntime::new(..).showing_progress()and gives itHANDSHAKE_TIMEOUTto reachis_established(). The pre-auth place is released at that point. - Connect. The runtime calls
KatanaConnector::connectfor every flow the core opens: the request, each mux sub-flow, each UDP association.connectadmits the user (a user the registry does not hold for that uid is refused withPermissionDenied"refused"), publishes the lease into the slot once, takes the source address (flow.source, else the peer) and builds aGate. A UDP flow becomes aFanOut. A TCP flow is routed withroute_target(destination, sniffed, source).Outbound::Blockor an audit hit refuses it withPermissionDenied"refused"; otherwise the chosen outbound is dialed and wrapped inMetered. - Relay.
drivepolls the runtime until it ends, until the lease is cancelled (until_retired), or until nothing has moved forPROGRESS_WATCHDOG. EveryTrafficdelta with a non-zero field resets the watchdog.
The detailed pages cover each step. Serving covers steps 1 to 4 and step 6. Admission and the connector and UDP fan-out cover step 5, and metering covers what Metered does with every byte.
Hysteria 2 nodes
Section titled “Hysteria 2 nodes”A Hysteria 2 node has no accept loop. TransportManager::start binds a std::net::UdpSocket and spawns run_hysteria, which hands the socket to the kernel:
pub async fn run_hysteria( inbound: Hy2Inbound<UserTag>, socket: std::net::UdpSocket, dispatcher: Arc<Dispatcher>, token: CancellationToken,)sequenceDiagram participant C as QUIC client participant H as Hy2Inbound participant KC as KatanaConnector participant AD as Admission participant O as Outbound C->>H: QUIC handshake, then auth H->>H: current Authenticator maps the credential to a UserTag C->>H: proxy stream, or UDP packets H->>KC: factory builds a connector for the client IP H->>KC: connect(flow) KC->>AD: admit(tag) AD-->>KC: counter and lease KC->>O: TCP only: route, audit, dial KC-->>H: Metered stream, or FanOut for UDP
The kernel listener authenticates each QUIC connection and runs its own runtime per proxy stream and, when UDP is enabled, one datagram runtime per connection. It calls the factory move |ip| KatanaConnector::new(dispatcher.clone(), Some(ip), None) for each of those runtimes, so every flow still passes admission, routing, audit and metering. There is no LeaseSlot and no katana-owned per-connection task. A retired user’s flows stop because their outbounds refuse to move: Gate::poll_open returns ConnectionAborted once the lease is cancelled, and new flows fail admission.
The kernel-side limits come from katana’s build_hysteria: max_connections: HY2_MAX_CONNECTIONS bounds QUIC connections, and circuit_permits, sized MAX_LIVE_CONNECTIONS_PER_NODE, bounds proxy streams and UDP sessions together. See Hysteria 2 server for the kernel side.
Compared with etemenanki-app
Section titled “Compared with etemenanki-app”The two programs share the kernel and the shape of the serve loop, and differ in everything around it.
| Stage | etemenanki-app | katana |
|---|---|---|
| Listener settings | The TOML file | The panel’s NodeInfo, plus listen_ip, certificates and [node.hysteria] from the TOML file |
| Protocols served | SOCKS, HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022, Hysteria 2, TUN | VMess, VLESS, Trojan, Shadowsocks, Shadowsocks 2022, Hysteria 2 |
| Accept loop | app/src/serve.rs → run_stream_inbound |
src/serve.rs → accept_loop |
| Task scope | spawn_scoped(token, fut) returns a JoinHandle |
spawn_scoped(&scope, fut) registers with the scope’s TaskTracker |
| Flow type | Flow<()> plus a FlowContext (inbound tag, source) |
Flow<UserTag>; the connector carries the source |
| Connector | AppConnector: route, then dial |
KatanaConnector: admit, route, audit, dial, meter |
| Route target | Destination, sniffed domain, network, inbound tag, source | Destination, sniffed domain, network, source; [[route.rule]] matches only domain suffix, CIDR, port, GeoSite and GeoIP |
| UDP association | FanOutLink: routes each packet |
FanOut: routes, audits and bills each packet |
| Ending one user’s connections | Not a concept: users are fixed for a generation | Lease cancellation through Admission |
| What a change costs | Any change to the config file builds a whole new generation; every connection drops | Per node: a panel user change swaps tables in place; a panel transport or protocol change, or a config-file edit to listen_ip, the certificate block, disable_sniffing, api.enable_vless, [node.hysteria], the route or the outbounds, rebuilds that node’s listener; any other [node.api] edit builds a new panel client and rebuilds only if the node info it reads changes the transport or protocol; a change to the panel identity respawns the node; other config fields apply without a rebuild |
| Hysteria 2 | run_hysteria_inbound, a connector without users |
run_hysteria, a connector that admits and meters |
The app’s side is described in app serving and generations and reload. The connection lifecycle both programs share is on connection lifecycle.
Control flow
Section titled “Control flow”The node task alternates between panel polls and static updates.
- Bootstrap (
NodeManager::bootstrap, onetry_bootstrapper attempt):forget_etags, so the panel answers in full;node_info;user_list; thenbring_up,set_curand, unlessdisable_get_ruleis set,refresh_rules. An attempt fails whennode_infoerrors or returns nothing, the port is0,user_listerrors or returns 304, orbring_upfails. The node then logsnode <id>: <reason>; retrying in <n>sand waits: 1 s after the first failure, doubling up to 60 s, and never longer than the poll period. AStaticUpdatethat arrives during the wait is applied; with nothing up it only stores what it carries, and aStaticUpdate::Configends the wait at once, because the edit may be the fix. The node token ends the wait or the attempt, andrunthen tears down whatever the attempt bound and sends its final reports. An empty user list is not a failure:bring_upbinds nothing and the node waits for the first poll that brings users. For a Hysteria 2 node with a non-zero[node.hysteria].port,node_infobuilds theNodeInfolocally and does not ask the panel. - Poll cycle (
poll_cycle(rebuild)), everyupdate_periodicseconds (at least 1): fetch node info and users, falling back tocuron an error or a 304; node info with port0also counts as a miss, so the node keeps the last node info and still applies the users;reconcile; refresh the rules unlessdisable_get_ruleis set;report_traffic;report_illegal. The timer passesrebuild = false. The first tick of the interval is consumed, so the first poll comes one full period after bootstrap. Withdisable_upload_trafficset,report_trafficsends nothing and discards residuals and writerless draining counters. The newV2board client’sreport_illegalis a no-op; only SSPanel receives audit hits, and it leaves out hits of local rules (rule id-1). - Reconcile ladder (
reconcile), first match wins:- No users: tear the listener down and commit an empty user set.
- A forced rebuild (
rebuild = true), no listener,NodeInfo::transport_eqfalse, orNodeInfo::protocol_eqfalse:rebuild, which tears down and then callsbring_up. user_set_differs, or the node speed limit changed:ProxyManager::refresh, which keeps the listener and the connections of every user who stays, rate-changed users included; departed and rebound users are retired.
- Static updates (
apply_static): aStaticUpdate::Configis taken whole or not at all. It first builds what the edit needs: a new panel client withPanelClient::newwhenpanel_typeor[node.api]changed, and a new router when the route changed. If either fails, the node logsconfig edit refused, keeping the running oneand changes nothing. Otherwise it stores the config, the router and the client. A change to the route,listen_ip, the certificate block,disable_sniffing,api.enable_vlessor[node.hysteria]setsrebuild. Whenrebuildis set or the client was replaced, and the node has bootstrapped,apply_staticrunspoll_cycle(rebuild)at once. A new client holds no ETags, so that poll reads the node and its users in full; the ladder then takes the narrowest step the answer needs, or the forced rebuild. An edit that changes only the client, for examplerule_list_pathorspeed_limit, therefore keeps every connection unless the new answer changes the transport or protocol.StaticUpdate::Outboundsswaps the pool, recompiles the router withrebuild_router, then rebuilds. Other fields are stored and read at their next use; a changedupdate_periodicrestarts the poll interval.
Node manager covers the ladder, and runtime and reload covers apply_reload. Panel clients covers what each fetch sends and parses.
Invariants
Section titled “Invariants”| Invariant | Mechanism | Pinned by |
|---|---|---|
| Everything that can fail is built before a listener binds, so a bad node never half-binds. | TransportManager::start builds the transport and protocol table, or the Hysteria listener, first and binds last; traffic.commit runs after the bind. |
Indirectly, by the builder refusals in tests/unit/inbound.rs, for example reality_rejected, vless_flow_rejected, tls_without_cert_file_errors, hysteria_requires_a_certificate |
| A user refresh never touches the listener or an unchanged user’s connection. | ProxyManager::refresh swaps ArcSwap<StreamProtocol> or calls set_authenticator; PreparedUsers keeps an unchanged user’s counter and leaves the user out of cancel_keys. |
unchanged_user_survives_user_refresh, repeated_user_refreshes_never_disturb_a_live_connection in tests/unit/e2e.rs |
| A departed or rebound user loses every flow and gets no new one. | refresh publishes the registry before the tables. Admission::commit swaps the registry and cancels the leases in cancel_keys under the leases lock that admit also takes. |
a_retired_user_stops_while_the_rest_keep_their_connections (tests/unit/e2e.rs); the_lease_reaches_the_connection_and_goes_with_the_user, a_credential_rebound_to_another_uid_is_refused (tests/unit/connector.rs); a_retired_users_connection_ends (tests/unit/serve.rs) |
| Only a registered credential is admitted, and it bills its own uid. | NodeTraffic::lookup filters on uid; tables carry UserTag, not counters. |
a_user_the_registry_does_not_know_is_refused, an_admitted_stream_is_billed_to_its_user (tests/unit/connector.rs) |
| A route or outbound pool edit drops every connection on the node. | The Router is fixed per generation. A route edit goes through poll_cycle(true), an outbound pool edit through apply_route_change, and both call rebuild. |
route_change_drops_connections (tests/unit/e2e.rs) |
| An edit that only changes the panel client takes effect without dropping a connection. | apply_static swaps the client and runs poll_cycle(false), and the ladder rebuilds only if the new answer changes the transport or protocol. |
a_client_edit_takes_effect_without_dropping_connections (tests/unit/runtime.rs) |
| A reload or a config edit that does not build changes nothing. | apply_reload builds every added node and every changed node’s panel client before it touches a running node; apply_static builds the client and router before it stores anything. |
a_reload_with_a_node_that_does_not_build_changes_nothing (tests/unit/runtime.rs) |
| A node that cannot come up keeps trying, and still stops when cancelled. | bootstrap retries with backoff and races every attempt and every wait against the node token. |
a_node_comes_up_once_the_panel_answers, a_node_whose_port_is_taken_comes_up_once_it_is_free, a_node_that_never_bootstraps_still_stops (tests/unit/e2e.rs) |
| A counter’s bytes are reported exactly once across a user change. | NodeTraffic::commit and snapshot hold users and draining together, in the same order, so a counter is in one place only. |
rate_change_drains_old_counter_and_reports_once, draining_counter_with_live_writer_is_retained, a_departed_users_late_bytes_are_still_reported (tests/unit/traffic.rs) |
| A report neither loses nor double-counts bytes on katana’s side. | commit_reported does a fetch_sub of the reported amounts only; a failed report puts counterless rows back with restore_residuals. |
commit_reported_preserves_concurrent, restored_residuals_are_retried (tests/unit/traffic.rs) |
| The speed limit is a property of bytes, not of chunk sizes. | TokenBucket::charge always takes n in full and goes into debt; Gate::poll_open waits until ready_at is None. |
a_chunk_larger_than_the_burst_is_still_limited, debt_holds_back_the_next_charge_too (tests/unit/traffic.rs); the_limit_holds_however_the_writes_are_sized (tests/unit/meter.rs) |
| A silent client cannot hold a pre-auth place past the handshake deadline. | drive wraps the whole handshake in tokio::time::timeout(HANDSHAKE_TIMEOUT, ..) and drops the permit on every exit. |
a_silent_client_is_dropped_at_the_handshake_deadline (tests/unit/serve.rs) |
| A connection that moves nothing is ended. | The PROGRESS_WATCHDOG sleep in drive, reset by every non-zero Traffic delta. |
a_connection_that_stops_moving_is_dropped (tests/unit/serve.rs) |
No parking_lot guard of the node is held across an .await. |
The guards are not Send, and the node future goes to tokio::spawn. |
The compiler |
| A teardown leaves no task of the generation behind. | Scope::shutdown cancels, closes the TaskTracker and waits; the Drop backstop cancels. |
No dedicated test; every test in tests/unit/e2e.rs cancels its node and awaits the task |
Failure paths and cancellation
Section titled “Failure paths and cancellation”| Failure | Effect |
|---|---|
At startup the config does not load, the outbound pool does not build, or there is no [[node]] |
run logs and returns ExitCode::FAILURE. |
| The config watcher cannot be set up | run logs config watcher disabled (no live reload) and keeps serving without reloads. |
| One node’s panel client or router fails to build at startup | spawn_node logs and returns None; the other nodes start. If none start, the process exits with a failure. |
A bootstrap attempt fails: node_info errors or returns nothing, the port is 0, user_list errors or returns 304, or bring_up fails |
bootstrap logs the reason and retries after the backoff (1 s, doubling to 60 s, capped at the poll period), or at once after a StaticUpdate::Config. Cancelling the node token ends the retries. |
A poll fetch fails or returns 304, or node info comes back with port 0 |
poll_cycle reuses the last applied node info or user list. The next tick is the retry. |
| A listener rebuild fails | The node has no listener. The next poll sees no transport and rebuilds again. |
| A user refresh fails to build its table | refresh logs and returns before publishing anything; the running tables and registry stay. |
A config edit’s route fails to compile, or its [node.api] does not build a panel client |
apply_static logs config edit refused, keeping the running one and stores nothing; the listener and its connections stay. |
| A rebuilt outbound pool leaves a node’s route uncompilable | rebuild_router keeps the old router, but apply_route_change still rebuilds the listener, so connections drop. |
| A reloaded config does not parse, has bad outbounds, or has a node that does not build | apply_reload is not called, or returns before applying anything; every node keeps running as it was. |
| A traffic report fails | Live and draining counters are untouched; counterless rows go back into residuals for the next cycle. |
| An audit report fails | The hits of that cycle are dropped and a warning is logged. |
accept fails with anything other than ConnectionAborted or Interrupted |
accept_loop backs off ACCEPT_ERROR_BACKOFF, or stops if the scope is cancelled meanwhile. |
Cancellation always travels downwards: root token → node token → the node’s own tear_down → Scope::shutdown → leases. drive futures are dropped where they wait, which drops the runtime, the client stream and every outbound it owns.
Limits
Section titled “Limits”| Constant | Value | Scope | Defined in |
|---|---|---|---|
MAX_LIVE_CONNECTIONS_PER_NODE |
65,536 | Accepted sockets per listener generation; also the Hysteria circuit_permits |
src/serve.rs |
MAX_TRANSPORT_STAGES_PER_NODE |
2,048 | Sockets in their TLS, WebSocket or HTTP/2 handshake; the accept loop waits when full | src/serve.rs |
TRANSPORT_HANDSHAKE_TIMEOUT |
10 s | How long a socket may hold its stage place | src/serve.rs |
PROGRESS_WATCHDOG |
RELAY_IDLE_TIMEOUT + 60 s = 360 s |
No bytes moved in either direction | src/serve.rs |
ACCEPT_ERROR_BACKOFF |
100 ms | After an accept error | src/serve.rs |
HANDSHAKE_TIMEOUT |
10 s | Until the core is established | kernel, protocols/src/core/mod.rs |
RELAY_IDLE_TIMEOUT |
300 s | The core’s own idle deadline | kernel, protocols/src/core/mod.rs |
MAX_PREAUTH_STREAMS_PER_NODE |
512 | Streams in their protocol handshake; excess streams are dropped | src/manager/proxy.rs |
HANDSHAKE_FAILURE_ALERT_PER_SEC |
10 | Failures per second above which the sampler logs a warning | src/manager/proxy.rs |
HY2_MAX_CONNECTIONS |
4,096 | QUIC connections per Hysteria 2 node | src/inbound.rs |
DRAIN_TIMEOUT, RELEASE_TIMEOUT |
3 s each | Waiting for a Hysteria endpoint to go idle, then for its port to come free | kernel, protocols/src/hysteria/server/inbound.rs |
MAX_SUBS |
64 | Outbound sub-links one UDP FanOut keeps open |
src/connector.rs |
MAX_RESOLVED_NAMES |
256 | Names one direct UDP association caches | src/outbound/freedom.rs |
| Static update channel | 16 | StaticUpdates queued per node |
src/runtime.rs → spawn_built |
| Reload debounce | 500 ms | Coalesces watcher events | src/runtime.rs → run |
| Poll period | update_periodic, at least 1 s |
poll_period |
src/manager/node.rs |
BOOTSTRAP_RETRY_MIN |
1 s | Wait after the first failed bootstrap attempt; doubles after each further failure | src/manager/node.rs |
BOOTSTRAP_RETRY_MAX |
60 s | Longest wait between bootstrap attempts; the poll period caps it when shorter | src/manager/node.rs |
| Panel request timeout | api.timeout, 5 s when 0 |
Every panel request | src/config.rs → timeout_secs |
The outbound client runtimes use a fixed buffer per protocol: HTTP_BUF, SOCKS_BUF and VLESS_BUF are 16 KiB, SS_BUF is 20 KiB, and VMESS_BUF and SS2022_BUF are 32 KiB (src/outbound/mod.rs).
Directorytests
- integration.rs one test target that spawns the real binary
Directoryintegration
- xray_interop.rs
- sniff.rs
- hysteria_interop.rs
Directorysupport
- mod.rs FakePanel, echo servers, certificates, reference client builds
Directoryunit
- e2e.rs NodeManager against fake UniProxy and SSPanel panels
- runtime.rs apply_reload against running nodes
- serve.rs
- connector.rs
- traffic.rs
- meter.rs
- rule.rs
- inbound.rs
- outbound.rs
Directoryapi
- newv2board.rs
- sspanel.rs
- Unit tests live under
tests/unit/but compile into the binary’s own test harness. Each tested source file ends with#[cfg(test)], a#[path]attribute pointing intotests/unit/, andmod tests;, so the tests can reach private items.src/main.rsincludestests/unit/e2e.rsthe same way, andsrc/runtime.rsincludestests/unit/runtime.rs. tests/unit/e2e.rsruns a realNodeManageragainst a hand-written UniProxy panel and drives it with the kernel’s own clients: VMess metering and reporting, user refresh, a proxy outbound, a route change, the bootstrap retries, and the Hysteria 2 node cases. The fake UniProxy panel can fail its first node config requests and can answer with ETags and304 Not Modified, as Xboard does; the bootstrap-retry tests use both. The file also holds a fake SSPanelmod_mupanel (fake_sspanel).tests/unit/runtime.rsdrivesapply_reloadagainst running nodes: which edits change a node’s identity (an_sspanel_node_is_its_panel_node_id,a_newv2board_node_is_also_the_type_it_asks_for,a_node_is_logged_by_its_panel_node_not_its_key), a newV2board type edit that respawns the node, a reload refused because a node does not build, and[node.api]edits applied in place: an SSPanelenable_vlessedit that rebuilds the listener, and arule_list_pathedit that keeps every connection.tests/integration.rsspawns the builtkatanabinary againstFakePaneland connects with the Xray and Hysteria reference clients.tests/support/mod.rsbuilds those clients withgo buildfrom theXray-coreandhysteriareference trees in the Etemenanki checkout. Whengoor a tree is missing, it printsSKIP: …and the test passes without running.- The tokio
test-utilfeature is enabled for tests, so the long deadlines intests/unit/serve.rs,tests/unit/traffic.rsandtests/unit/meter.rsrun on a paused clock (#[tokio::test(start_paused = true)]).
See testing for how to run the suites.