Design principles
Source files: 70 · checked against Etemenanki 596916d · katana v3.0.1
Etemenanki/concepts/src/lib.rsEtemenanki/concepts/src/core.rsEtemenanki/concepts/src/runtime.rsEtemenanki/concepts/src/buffer.rsEtemenanki/concepts/src/wake.rsEtemenanki/concepts/src/client.rsEtemenanki/concepts/tests/runtime.rsEtemenanki/concepts/tests/client.rsEtemenanki/environment/src/lib.rsEtemenanki/protocols/src/lib.rsEtemenanki/protocols/src/error.rsEtemenanki/protocols/src/helpers/parse.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/core/harness.rsEtemenanki/protocols/src/sniff/mod.rsEtemenanki/protocols/src/vmess/core.rsEtemenanki/protocols/src/vmess/accounts.rsEtemenanki/protocols/src/ss_2022/core.rsEtemenanki/protocols/src/ss_legacy/core.rsEtemenanki/protocols/src/trojan/core.rsEtemenanki/protocols/src/vless/core.rsEtemenanki/protocols/src/http/core.rsEtemenanki/protocols/src/tun/udp.rsEtemenanki/protocols/src/tun/config.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/protocols/src/wireguard/device.rsEtemenanki/protocols/src/transports/grpc/stream.rsEtemenanki/protocols/src/hysteria/connection.rsEtemenanki/protocols/src/hysteria/server/config.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/hysteria/server/datagrams.rsEtemenanki/protocols/tests/unit/core/mod.rsEtemenanki/protocols/tests/unit/trojan/core.rsEtemenanki/protocols/tests/unit/vmess/core.rsEtemenanki/protocols/tests/unit/vmess/protocol.rsEtemenanki/protocols/tests/unit/ss_2022/core.rsEtemenanki/protocols/tests/unit/mux/demux.rsEtemenanki/protocols/tests/unit/wireguard/device.rsEtemenanki/protocols/tests/pipeline/wireguard.rsEtemenanki/protocols/tests/unit/mux/frame.rsEtemenanki/protocols/tests/unit/hysteria/protocol.rsEtemenanki/protocols/tests/unit/hysteria/server/datagrams.rsEtemenanki/protocols/tests/unit/transports/grpc_liveness.rsEtemenanki/app/src/config.rsEtemenanki/app/src/transport.rsEtemenanki/app/src/inbound/mod.rsEtemenanki/app/src/inbound/tun.rsEtemenanki/app/src/outbound/mod.rsEtemenanki/app/src/main.rsEtemenanki/app/src/instance.rsEtemenanki/app/src/serve.rsEtemenanki/app/src/outbound/udp_fanout.rsEtemenanki/app/tests/unit/config.rsEtemenanki/app/tests/unit/transport.rsEtemenanki/app/tests/unit/inbound.rsEtemenanki/app/tests/integration/e2e_hysteria.rsEtemenanki/app/tests/integration/e2e_hysteria_inbound.rskatana/src/serve.rskatana/src/connector.rskatana/src/meter.rskatana/src/traffic.rskatana/src/config.rskatana/src/manager/proxy.rskatana/src/manager/transport.rskatana/src/manager/node.rskatana/src/runtime.rskatana/tests/unit/meter.rskatana/tests/unit/traffic.rskatana/tests/unit/runtime.rs
This page collects the rules the code base is built on. Each rule comes with the reason for it, the type or check that enforces it, and the tests that pin it. The last section lists the review rules that maintainers apply to every change, stated as design rules.
Read it before changing a protocol core, the per-connection runtime, the application’s build and reload path, or katana’s metering. Most of these rules are carried by the shape of the types rather than by convention. A change that seems to need a way around one of them usually belongs in a different layer.
At a glance
Section titled “At a glance”| Principle | Enforced by | Pinned by |
|---|---|---|
| Protocol cores are sans-I/O | ProxyCoreDecode::handle takes an Event and an Effects sink and nothing else; clocks are constructor arguments |
core_is_driven_without_any_io, codec_is_driven_without_any_io, every CoreHarness test |
| One task per connection | ProxyServerRuntime owns the transport, the core and every outbound; ProxyClientRuntime is an AsyncRead + AsyncWrite |
server_runtime_relays_through_a_client_runtime_in_one_task, stalled_outbound_holds_uplink_but_not_other_downlink |
| Fixed buffers, explicit caps | three boxed [u8; BUF_SIZE] arrays per connection; semaphores and table-size checks with try_acquire |
frame_larger_than_the_buffer_is_an_error, unknown_or_excess_sessions_are_declined_with_end |
| Errors are events | RuntimeError has no outbound variant; outbound failures arrive as Event::ConnectFailed, OutboundError and SendFailed |
connect_failure_reaches_the_core_as_an_event, a_refused_datagram_send_keeps_the_key_alive |
| Fail-closed configuration | #[serde(deny_unknown_fields)] on every config struct; behaviour-selecting strings matched against explicit lists in the builders |
a_mistyped_key_is_rejected_rather_than_ignored, an_unknown_security_is_rejected_on_every_network |
| Build, then swap | instance::build binds nothing; Instance::reload touches the old generation only after build succeeds; katana’s apply_reload builds every added node before it applies anything; katana’s NodeTraffic::prepare / commit |
a_reload_rebinds_the_udp_port, katana a_reload_with_a_node_that_does_not_build_changes_nothing, katana rate_change_drains_old_counter_and_reports_once |
| Deterministic time in tests | deadlines are events; the runtime arms tokio::time timers; katana’s TokenBucket reads tokio::time::Instant |
deadline_is_armed_against_the_tokio_clock, a_chunk_larger_than_the_burst_is_still_limited |
| Panic-free parsing | crate-level #![deny(clippy::…)] in etemenanki-protocols and etemenanki-environment; helpers::parse::{take, take_array, need_more} |
varint_above_62_bits_is_refused_not_panicked, udp_truncated_in_the_header_is_refused_and_never_panics |
Protocol cores are sans-I/O
Section titled “Protocol cores are sans-I/O”A protocol’s server side is a state machine from bytes and events to bytes and effects. It does not own a socket, a timer or a task Context. The module documentation of concepts/src/core.rs states the rule in one line: “Neither core touches a socket, a clock or a Context”.
Why. A proxy protocol is mostly parsing, cryptography and bookkeeping. With no I/O inside, a unit test calls the core with a hand-built byte slice and checks exactly what it asked for. The same core then runs unchanged over TCP, a Unix socket, a WebSocket, an HTTP/2 stream or a QUIC stream. Every scheduling decision (what to read, when to stop reading, when to give up) lives in one place, the runtime, instead of being repeated in every protocol module.
The contract
Section titled “The contract”concepts/src/core.rs → ProxyCoreDecode is the whole interface between a server protocol and the rest of the system:
pub trait ProxyCoreDecode { type Key: Copy + Ord + Send + Sync + 'static; type Target; type Error; type TransportAddr: Clone + Send + Sync + 'static;
const STAGING_RESERVE: usize; const MAX_DATAGRAM: usize = 4096;
fn handle( &mut self, event: Event<'_, Self>, effects: &mut Effects<'_, Self>, ) -> Result<usize, Self::Error>;
fn held(&self) -> &[u8] { … }}A core receives one Event per call and returns how many leading bytes of it the core consumed. It answers by pushing Effects into the sink and by staging reply bytes toward the transport:
| Kind | Variants |
|---|---|
| Events from the transport | Transport, TransportDatagram, TransportSendFailed, TransportEof |
| Events from an outbound | Outbound, Datagram, SendFailed, OutboundEof, Connected, ConnectFailed, OutboundError |
| Events from the timer | Deadline |
| Effects that move bytes | Forward, SendTo (a range of the event’s slice); ForwardHeld, SendToHeld (a range of the core’s held buffer) |
| Effects on outbounds | Open, Shutdown, Close |
| Effects on the connection | ShutdownTransport, SetDeadline(Option<Duration>), Finish |
The client side (ProxyCoreEncodeHandshake, ProxyCoreEncode, ProxyCoreEncodeDatagram) follows the same rule: a codec seals plaintext into a Staging area and opens wire frames in place, with no I/O.
Hidden inputs are constructor arguments
Section titled “Hidden inputs are constructor arguments”The core module names the three inputs that would otherwise make a core non-deterministic: “Randomness, wall-clock time and shared account state are constructor arguments of the concrete core, never something the runtime injects”. The two cores that check timestamps take the clock as a plain function pointer:
impl<T> VMessCore<T> { pub fn new( validator: Arc<AccountValidator<T>>, now: fn() -> i64, sniff: bool, source: Option<IpAddr>, ) -> Self}impl<T> Ss2022Core<T> { pub fn new( config: Arc<Ss2022ServerConfig<T>>, validator: Option<Arc<Validator<T>>>, sniff: bool, source: Option<IpAddr>, now: fn() -> u64, ) -> Self
pub fn with_system_clock( config: Arc<Ss2022ServerConfig<T>>, validator: Option<Arc<Validator<T>>>, sniff: bool, source: Option<IpAddr>, ) -> Self}app/src/serve.rs passes the real clock (vmess::aead::now_unix, or Ss2022Core::with_system_clock). A test passes a constant: eih_selects_the_user_and_a_stale_timestamp_is_refused in protocols/tests/unit/ss_2022/core.rs builds a core with the clock || 1_000 and checks that the request is refused as stale.
Shared account state arrives the same way, and it may keep time of its own. The VMess AccountValidator passed to VMessCore::new expires its replay window against std::time::Instant, so the injected now decides the timestamp check but not how long an auth ID counts as seen.
One deadline, armed by effect
Section titled “One deadline, armed by effect”Timeouts follow the same rule. A core does not measure its own timeouts: it arms the runtime’s single timer with Effect::SetDeadline and reacts to Event::Deadline. The contract in concepts/src/core.rs says the current time “comes from a clock passed to the constructor, never from Instant::now() inside handle”, and that a protocol with several timeouts keeps its own ordered map of expiries, arms the earliest and re-arms on every Deadline.
protocols/src/core/mod.rs → Timing is the shared helper for the common case. The connection’s Phase decides what the one deadline means:
stateDiagram-v2 [*] --> Handshake Handshake --> Sniff: request parsed, IP destination, sniffing on Handshake --> Relay: request parsed Sniff --> Relay: domain found, limit reached or deadline Relay --> Closing: one side closed Relay --> Closing: idle deadline, Finish pushed Closing --> [*]
| Phase | Deadline armed | Constant | Value |
|---|---|---|---|
Handshake |
once, on the first byte event (Timing::touch) |
HANDSHAKE_TIMEOUT |
10 s |
Sniff |
on entering the phase (Timing::enter) |
SNIFF_TIMEOUT (protocols/src/sniff/mod.rs) |
300 ms |
Relay, Closing |
again on every byte event, so it measures idleness, not lifetime | RELAY_IDLE_TIMEOUT |
300 s |
When the idle deadline passes, Timing::expired pushes Effect::Finish itself and returns Expired::Idle. For Handshake and Sniff it only reports which deadline passed, and the core decides.
Testing a core with CoreHarness
Section titled “Testing a core with CoreHarness”protocols/src/core/harness.rs → CoreHarness stands in for the runtime. It owns an EffectList and a WriteBuffer<HARNESS_STAGING> (HARNESS_STAGING = 64 KiB), and does no I/O:
pub struct CoreHarness<C: ProxyCoreDecode> { pub core: C, // effects, staging and the datagram packet list are private}
impl<C: ProxyCoreDecode> CoreHarness<C> { pub fn new(core: C) -> Self pub fn over_datagrams(core: C) -> Self pub fn event(&mut self, event: Event<'_, C>) -> Result<(usize, Vec<Effect<C>>), C::Error> pub fn transport(&mut self, data: &mut [u8]) -> Result<(usize, Vec<Effect<C>>), C::Error> pub fn feed(&mut self, data: &mut [u8]) -> Result<(usize, Vec<Effect<C>>), C::Error> pub fn outbound( &mut self, key: C::Key, data: &mut [u8], ) -> Result<(usize, Vec<Effect<C>>), C::Error> pub fn staged(&mut self) -> Vec<u8> pub fn staged_packets(&mut self) -> Vec<(Vec<u8>, C::TransportAddr)> pub fn held(&self, range: std::ops::Range<usize>) -> Vec<u8>}feed behaves like the runtime: it offers the unconsumed tail again while the core makes progress, and rebases Forward and SendTo ranges onto the caller’s slice. A timeout test needs no clock, because the deadline is only an event:
#[test]fn handshake_deadline_fails_the_connection() { let mut h = core(false); let mut partial = vec![0u8; 10]; h.transport(&mut partial).unwrap(); let err = h.event(Event::Deadline).unwrap_err(); assert_eq!(err.kind(), io::ErrorKind::TimedOut); let _ = Duration::ZERO;}The sans-I/O shape also makes the subtle parts of the contract testable. The core documentation says “Never decrypt what you do not consume”: a core that decrypts a length header in place must remember that it did, because the same bytes come back on the next call. a_chunk_split_across_reads_decodes_its_header_once in protocols/tests/unit/vmess/core.rs cuts a VMess request ten bytes short, feeds the rest, and checks that the payload still comes out intact.
| Test | File | What it pins |
|---|---|---|
core_is_driven_without_any_io |
concepts/tests/runtime.rs |
A toy mux core is driven with a bare EffectList and WriteBuffer; a split frame stays unconsumed. |
codec_is_driven_without_any_io |
concepts/tests/client.rs |
The client codec contract works the same way. |
timing_arms_handshake_once_then_idle_per_byte_event |
protocols/tests/unit/core/mod.rs |
Timing arms the handshake deadline once, then the idle deadline on every byte event. |
a_chunk_split_across_reads_decodes_its_header_once |
protocols/tests/unit/vmess/core.rs |
The “never decrypt twice” rule. |
eih_selects_the_user_and_a_stale_timestamp_is_refused |
protocols/tests/unit/ss_2022/core.rs |
The injected clock decides the timestamp check. |
handshake_deadline_fails_the_connection |
protocols/tests/unit/trojan/core.rs |
A deadline delivered by hand fails a half-finished handshake with TimedOut. |
One task per connection
Section titled “One task per connection”One proxied connection is exactly one future: concepts/src/runtime.rs → ProxyServerRuntime. It owns the client’s transport, the core, the connector and every outbound the core opens, and it moves every byte itself. There is no task per direction and no channel between the halves.
pub struct ProxyServerRuntime<const BUF_SIZE: usize, Core, Trans, Conn, Mode = ProxyRunsQuiet>where Core: ProxyCoreDecode, Conn: Connector<Core::Target>,{ … }
impl<const BUF_SIZE: usize, Core, T, Conn> ProxyServerRuntime<BUF_SIZE, Core, StreamTransport<T>, Conn, ProxyRunsQuiet>where Core: ProxyCoreDecode, T: AsyncRead + AsyncWrite + Unpin, Conn: Connector<Core::Target>,{ pub fn new(transport: T, core: Core, connector: Conn) -> Self}With the default ProxyRunsQuiet marker the runtime is a Future whose output is Result<Traffic, RuntimeError<Core::Error>>. showing_progress() turns it into a Stream of Traffic deltas (the ProxyShowsProgress marker) for a caller that accounts traffic as it moves. over_datagrams builds the same runtime over a DatagramLink transport instead of a byte stream.
An upstream proxy does not break the rule. concepts/src/client.rs → ProxyClientRuntime implements AsyncRead + AsyncWrite on its plaintext side and owns the wire side, so the server runtime polls it like any other outbound stream. Its module documentation: the server core “forwards plaintext into it with poll_write, reads plaintext back with poll_read, and nothing in between is a task or a channel”.
flowchart LR
client["client wire"]
subgraph task["one task: ProxyServerRuntime"]
up["up: ReadBuffer"]
core["ProxyCoreDecode"]
staging["staging: WriteBuffer"]
scratch["scratch buffer"]
out1["outbound stream"]
out2["ProxyClientRuntime"]
end
dest["destination"]
upstream["upstream proxy"]
client --> up --> core
core -- "Forward" --> out1 --> dest
core -- "Forward" --> out2 --> upstream
out1 -- "read" --> scratch --> core
core -- "stage" --> staging --> client
Why. With a task per direction and channels between them, each channel is a queue that someone has to bound, each task is a lifetime that someone has to end, and each error has to be carried across a channel to whoever decides what to do about it. With one task:
- Backpressure follows from the structure. The runtime reads from a side only when there is room for the result. When a destination stops accepting writes, the forward waits, and the transport is not read.
- Cancellation is by drop. Dropping the future drops the transport, every outbound and every buffer. Nothing is left running that would need to be found and stopped.
- Ordering is total. Effects apply in the order the core pushed them, so a core can reason about “open, then forward, then close” without races.
How backpressure works
Section titled “How backpressure works”The # Backpressure section of concepts/src/runtime.rs states the rules the scheduler follows:
- Effects apply strictly in order, and a forward that cannot complete holds everything behind it. While a forward of transport bytes waits, the transport is not read. For a multiplexing core, a stalled sub-flow therefore holds its siblings’ uplink until it drains; downlink from other outbounds keeps flowing.
- A stream outbound is read only while staging has
STAGING_RESERVE + 1bytes free. A datagram outbound is read only while staging hasSTAGING_RESERVEplus the datagram limit free. A client that stops reading stops the downlink at the outbound sockets, and a packet is never read into less room than it may need. - A held-range effect (
ForwardHeld,SendToHeld) pins the core’s held buffer. No byte event is delivered until it has been applied.
Only outbounds whose waker fired are polled. concepts/src/wake.rs → ReadyQueue and KeyWaker record which key became ready, so a connection with many outbounds does not poll all of them on every wake-up. In the Future mode, poll performs at most WORK_BUDGET (64) units of work, then wakes itself and returns Pending, so a connection whose sockets are always ready cannot monopolise a worker thread.
Cancellation by drop
Section titled “Cancellation by drop”The application ends connections by dropping their futures. app/src/serve.rs → spawn_scoped ties every spawned connection to its generation’s CancellationToken:
pub fn spawn_scoped<F>(token: CancellationToken, fut: F) -> JoinHandle<()>where F: Future + Send + 'static, F::Output: Send,It runs tokio::select! between token.cancelled() and the future. When the generation is cancelled, the connection future is dropped along with everything it owns. katana’s src/serve.rs → spawn_scoped(scope: &Scope, fut: F) does the same under a Scope, which adds a TaskTracker so that Scope::shutdown returns only after every task under it is gone.
Two client-side carriers have a driver task of their own. A gRPC client stream (protocols/src/transports/grpc/stream.rs → GrpcStream::connect) sits on an HTTP/2 connection whose connection future has to be polled apart from the stream. A Hysteria 2 client connection (protocols/src/hysteria/connection.rs) is shared by every circuit routed to it and runs its HTTP/3 driver, plus a datagram pump when the server said it relays UDP. The flows on top still run one task each. Both files hold these tasks in an AbortOnDropHandle, so dropping the owner also stops the driver.
katana meters inside the same task
Section titled “katana meters inside the same task”katana keeps this model for accounting. src/connector.rs wraps each stream outbound it hands to the runtime in src/meter.rs → Metered<S>. A Gate bills bytes to the user’s UserCounter and holds the flow back while the user’s token bucket is in debt:
impl Gate { pub fn new(counter: Arc<UserCounter>, retired: CancellationToken) -> Self pub fn poll_open(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>> pub fn sent(&self, n: usize) pub fn received(&self, n: usize)}Metered::poll_read and poll_write call poll_open first. A speed limit is therefore ordinary backpressure: they return Pending until the debt is paid, and the runtime stops reading the other side. There is no relay loop to count in and no extra task per user. The same gate also ends the flow with ConnectionAborted once the user’s lease token is cancelled.
The debt comes from src/traffic.rs → TokenBucket. Gate::sent and Gate::received call TokenBucket::charge with the byte count after each transfer, and its documentation states the rule: “A charge is always taken in full, even when it is larger than what the bucket holds: the balance goes negative”. Before the next transfer, poll_open asks ready_at for the instant the debt is repaid and sleeps until then. The burst is one second of rate, and rate == 0 means unlimited.
| Test | File | What it pins |
|---|---|---|
stalled_outbound_holds_uplink_but_not_other_downlink |
concepts/tests/runtime.rs |
A stalled forward stops the uplink; another key’s downlink still reaches the client. |
datagram_outbound_is_never_truncated_by_staging_backpressure |
concepts/tests/runtime.rs |
A datagram is read only into enough room. |
server_runtime_relays_through_a_client_runtime_in_one_task |
concepts/tests/client.rs |
A server runtime relays through a client runtime without a second task. |
backpressure_from_the_wire_reaches_the_writer |
concepts/tests/client.rs |
A slow upstream makes the plaintext writer wait. |
the_limit_is_shared_by_both_directions |
katana tests/unit/meter.rs |
One bucket gates both directions of a flow. |
retiring_the_user_wakes_a_parked_read |
katana tests/unit/meter.rs |
Retiring a user ends a flow even while it waits on a silent peer. |
debt_holds_back_the_next_charge_too |
katana tests/unit/traffic.rs |
A 25,000-byte charge against a 10,000 B/s bucket leaves 1.5 s of debt that the next charge waits out. |
Fixed buffers and explicit caps
Section titled “Fixed buffers and explicit caps”Every per-connection resource has a size or a count that is fixed in the code and visible in review.
Buffers never grow
Section titled “Buffers never grow”A ProxyServerRuntime allocates three boxed arrays of BUF_SIZE bytes, once: up: ReadBuffer<BUF_SIZE> for transport reads, staging: WriteBuffer<BUF_SIZE> toward the transport, and scratch: Box<[u8; BUF_SIZE]> for outbound reads. Its documentation says the connection allocates “exactly those plus one small Arc per outbound opened” (the outbound’s waker). concepts/src/buffer.rs puts it directly: the “buffers never grow; backpressure comes from them being full.”
The consequences are enforced:
- The runtime constructor asserts
BUF_SIZE > Core::STAGING_RESERVE(“BUF_SIZE must exceed the core’s STAGING_RESERVE or nothing is ever read”). - A frame that does not fit is an error, not a reallocation. When the unparsed region fills the buffer and the core still wants more, the runtime fails the connection with
RuntimeError::FrameTooLarge. - Each core declares its size as an associated
BUF_SIZE, and the application instantiates the runtime with it, for exampledrive::<{ TrojanCore::<()>::BUF_SIZE }, _, _>.
| Core | BUF_SIZE |
Reason given in the source |
|---|---|---|
PassthroughCore |
8 KiB | Not stated in the source. The core relays verbatim; the TUN inbound uses it for TCP flows. |
TunUdpCore |
8 KiB | A datagram plus the runtime’s reserve. |
Hy2StreamCore |
8 KiB | A request with the longest address and padding. |
TrojanCore |
16 KiB | A UDP packet frame: the header plus at most 8 KiB. |
VlessCore |
16 KiB | A UDP frame of at most MAX_DATAGRAM bytes behind its length. |
Hy2UdpCore |
16 KiB | A reassembled packet plus its fragments’ headers. |
ShadowsocksCore |
20 KiB | One chunk of MAX_PAYLOAD plus its overhead, with room to spare. |
VMessCore |
32 KiB | Not stated in the source. |
Ss2022Core |
32 KiB | A request chunk with full padding, or a record of the size peers send. |
HttpCore |
MAX_HEAD = 64 KiB |
A whole HTTP request head. |
Counts are capped with permits
Section titled “Counts are capped with permits”A cap is a semaphore or a table-size check. Most are applied with try_acquire_owned or a length check: when the cap is reached, the new item is refused at once instead of waiting in a queue that would itself need a bound. The comment in katana’s accept loop gives the reason for refusing rather than queueing: waiting there “would stop accepting altogether and let the listen backlog absorb the overload instead, which is how a saturated node turns into a silent one.”
| Cap | Where | Value | Over the cap |
|---|---|---|---|
MAX_LIVE_CONNECTIONS_PER_INBOUND |
app/src/serve.rs |
65,536 per TCP or Unix-socket inbound | The accepted socket is dropped. The owned permit travels with the socket into every stream it carries. |
MAX_LIVE_CONNECTIONS_PER_NODE |
katana src/serve.rs |
65,536 per listener | The accepted socket is dropped. The permit is held for the whole connection, every stream it carries included. |
MAX_SESSIONS |
protocols/src/mux/demux.rs |
256 per mux carrier | The new sub-flow is answered with End; the carrier survives. |
DEFAULT_MAX_CONNECTIONS |
protocols/src/hysteria/server/config.rs |
4,096 per listener (max_connections) |
The incoming QUIC connection is refused. |
DEFAULT_MAX_CIRCUITS |
protocols/src/hysteria/server/config.rs |
65,536 per listener (max_circuits) |
A new proxy stream is reset with H3_REQUEST_REJECTED; a new UDP association is dropped. |
MAX_SESSIONS |
protocols/src/hysteria/server/datagrams.rs |
256 per QUIC connection | A packet that would open another UDP association is dropped. |
MAX_UDP_SESSIONS |
protocols/src/hysteria/connection.rs |
256 per Hysteria 2 client connection | Opening another UDP session fails with WouldBlock. |
MAX_SUBS |
app/src/outbound/udp_fanout.rs |
64 sub-links (one per routed outbound) per UDP association | The sub-link least recently sent to is closed first. |
DEFAULT_MAX_FLOWS |
protocols/src/tun/config.rs |
65,536 per device (max_flows) |
The new flow is dropped. |
The live-connection caps are guardrails, not quotas. The comment on MAX_LIVE_CONNECTIONS_PER_INBOUND explains the sizing: far above normal traffic, and still below the process file-descriptor limit once each session’s outbound socket is counted. The Hysteria 2 inbound’s two configurable caps must be at least 1: max_connections = 0 and max_circuits = 0 are refused when the configuration is built (“must be at least 1”). The TUN inbound’s max_flows has no such check.
| Test | File | What it pins |
|---|---|---|
frame_larger_than_the_buffer_is_an_error |
concepts/tests/runtime.rs |
A 500-byte frame into a 128-byte runtime ends with FrameTooLarge. |
unknown_or_excess_sessions_are_declined_with_end |
protocols/tests/unit/mux/demux.rs |
The mux session cap declines, it does not tear down. |
the_circuit_limit_refuses_new_sessions |
protocols/tests/unit/hysteria/server/datagrams.rs |
A UDP association over the circuit limit is refused. |
stream_permits_are_released_when_a_circuit_ends |
app/tests/integration/e2e_hysteria.rs |
The Hysteria 2 client’s stream permits come back when a circuit ends, so 20 circuits in a row pass a limit of 4. Builds the upstream Hysteria server with go and returns early when it cannot. |
a_zero_limit_is_refused |
app/tests/unit/inbound.rs |
max_connections and max_circuits must be at least 1. |
Errors are events
Section titled “Errors are events”A connection has one client and possibly many outbounds. If an outbound’s failure were a runtime error, one dead destination would end every sub-flow on a mux carrier, and the core could never answer the client with an HTTP 502 or a mux End. So the runtime reports outbound failures to the core as events, and the core decides.
concepts/src/runtime.rs → RuntimeError has no variant for an outbound. Its documentation: “Outbound failures are not here: they reach the core as events and it decides.”
pub enum RuntimeError<E> { Transport(io::Error), Core(E), UnknownKey, DuplicateKey, RangeOutOfBounds, BadConsume, FrameTooLarge, WrongLinkKind, StagedWithoutPeer,}Transport is the client’s side failing, Core is the protocol refusing the client, and FrameTooLarge is a frame that does not fit the read buffer. The other variants are a core breaking its contract, for example forwarding a range outside the event’s slice (RangeOutOfBounds) or staging bytes toward a datagram transport without a peer (StagedWithoutPeer).
| Failure | Arrives as | Key afterwards |
|---|---|---|
The dial for Effect::Open fails |
Event::ConnectFailed { key, error } |
Gone |
| An outbound read or write fails | Event::OutboundError { key, error } |
Gone |
| A datagram outbound refuses one packet | Event::SendFailed { key, to, error } |
Still live |
| A datagram transport refuses one packet to a peer | Event::TransportSendFailed { to, error } |
Not applicable; the transport keeps serving other peers |
sequenceDiagram participant C as client participant R as ProxyServerRuntime participant K as HttpCore participant D as connector C->>R: CONNECT request R->>K: Event::Transport K->>R: Effect::Open R->>D: connect(target) D-->>R: io::Error R->>K: Event::ConnectFailed K->>R: stage RESP_502, ShutdownTransport, Finish R->>C: HTTP 502, then close
The client side makes this precise for upstream proxies. The connect future of ProxyClientConnector resolves only after the dial, the codec’s handshake and a flush of the handshake bytes (poll_connected). A server core’s Event::Connected therefore means the upstream flow is really open, and an upstream that refuses the handshake becomes ConnectFailed, not a relay that dies on its first byte.
Inside etemenanki-protocols, parsers classify failures as protocols/src/error.rs → ProtocolError and convert them to io::Error with a fixed kind:
ProtocolError variant |
io::ErrorKind |
|---|---|
Truncated |
UnexpectedEof |
Unauthenticated |
PermissionDenied |
Malformed, Overflow, Unsupported |
InvalidData |
Crypto, Other |
Other |
Io(e) |
the kind of e |
need_more relies on this mapping: it turns UnexpectedEof into Ok(None) (“wait for more bytes”) and keeps every other error.
| Test | File |
|---|---|
connect_failure_reaches_the_core_as_an_event |
concepts/tests/runtime.rs |
a_refused_datagram_send_keeps_the_key_alive |
concepts/tests/runtime.rs |
a_refused_transport_packet_is_reported_not_fatal |
concepts/tests/runtime.rs |
upstream_dial_failure_is_connect_failed_not_connected |
concepts/tests/client.rs |
refused_upstream_handshake_is_connect_failed_not_connected |
concepts/tests/client.rs |
Fail-closed configuration
Section titled “Fail-closed configuration”A configuration mistake stops the program with an error. It never falls back to a default that means something else. The comment on InboundConfig::listen in app/src/config.rs names the failure this prevents: a listen key that fails to deserialise “silently falling back to 0.0.0.0 and putting a proxy that was meant to be local on the network”.
Unknown keys are errors
Section titled “Unknown keys are errors”Every struct in app/src/config.rs carries #[serde(deny_unknown_fields)]: the top-level Config, every section, and every per-protocol settings struct. Per-protocol settings stay an opaque toml::Value until the builder knows the protocol; app/src/inbound/mod.rs → parse_settings then deserialises them into a struct that also denies unknown fields. Every struct in katana’s src/config.rs carries the same attribute.
$ etemenanki-app --test -c config.tomlERROR etemenanki_app: configuration invalid: TOML parse error at line 4, column 1 |4 | lisen = "0.0.0.0" | ^^^^^unknown field `lisen`, expected one of `tag`, `protocol`, `listen`, `port`, `stream`, `address_family`, `sniffing`, `settings`A missing listen binds 127.0.0.1, not every interface. The comment gives the reasoning: deny_unknown_fields catches the typo, “and the loopback default catches whatever it does not: a server that wants every interface says so.”
Unknown values are errors
Section titled “Unknown values are errors”Strings that select behaviour are matched against an explicit list in the builders, and anything else is an error. app/src/transport.rs shows the pattern:
pub fn tls_layer(network: &str, security: Option<&str>, ctx: &str) -> io::Result<bool>pub fn resolve_stream(stream: &StreamConfig, ctx: &str) -> io::Result<StreamShape>pub fn reject_stream_settings(stream: &StreamConfig, proto: &str, ctx: &str) -> io::Result<()>tls_layer validates security against the network instead of comparing it to "tls" with plaintext as the fallback. Its comment explains why this matters for a proxy in particular: when a plaintext transport is built by mistake, the only symptom is a failed handshake after the credential has already crossed the network in the clear. reject_stream_settings covers the related case of a [stream] block on a protocol that never reads it (SOCKS, Shadowsocks, Hysteria 2 and TUN inbounds, and any inbound on a Unix socket; freedom, blackhole, hysteria2 and wireguard outbounds), which would otherwise be discarded without a word.
Each of these was checked with etemenanki-app --test:
| Mistake | Error |
|---|---|
lisen = "0.0.0.0" |
unknown field `lisen`, expected one of … |
protocol = "sock" |
inbound in: unknown protocol "sock" |
security = "tsl" |
inbound in: unknown stream security "tsl" (expected "tls" or "none") |
network = "tcp" with security = "tls" |
inbound in: security = "tls" is not valid with network = "tcp"; use network = "tls" for TLS over plain TCP (…) |
network = "ws" on a SOCKS inbound |
inbound in: protocol socks does not support stream network "ws" |
udpp = 1 in SOCKS settings |
inbound in: invalid settings: unknown field `udpp`, expected one of `auth`, `accounts`, `udp`, `udp_bind` |
network = "tpc" in a route rule |
invalid rule network "tpc" (expected "tcp" or "udp") |
domain_regex = ["("] |
invalid domain regex "(": regex parse error: … |
port = ["2000-1000"] |
invalid port spec: "2000-1000" has a lower bound above its upper bound |
cidr = ["10.0.0.0/33"] |
invalid cidr "10.0.0.0/33": … |
outbound = "direkt" in a route rule |
route references unknown outbound tag: direkt |
--test runs config::load and then the same instance::build that a real start runs, so it refuses exactly what a start would refuse, without binding anything.
| Test | File |
|---|---|
a_mistyped_key_is_rejected_rather_than_ignored |
app/tests/unit/config.rs |
an_inbound_defaults_to_loopback |
app/tests/unit/config.rs |
an_unknown_security_is_rejected_on_every_network |
app/tests/unit/transport.rs |
tcp_with_tls_is_rejected_and_names_the_fix |
app/tests/unit/transport.rs |
a_transport_the_protocol_cannot_honour_is_rejected |
app/tests/unit/transport.rs |
an_unknown_setting_is_refused |
app/tests/unit/inbound.rs |
Build, then swap
Section titled “Build, then swap”A new configuration is parsed, validated and fully constructed before anything that serves traffic is touched. app/src/instance.rs → build produces a Built value (the router, every inbound with its BindSpec, the balancers and the resolver) and binds nothing:
pub fn build(cfg: &Config) -> io::Result<Built>Instance::reload uses it as a gate:
flowchart TB
read["read file"] --> same{"bytes unchanged?"}
same -- "yes" --> stop["return"]
same -- "no" --> parse["config::parse_bytes"]
parse -- "error" --> keep["log, keep old generation"]
parse --> build["instance::build"]
build -- "error" --> keep
build --> cancel["cancel old generation, await its accept loops"]
cancel --> spawn["spawn_generation(built, false)"]
A configuration that fails to parse or build is logged, and the old generation keeps running untouched. Because each inbound’s BindSpec (TCP, UDP, Unix path or TUN device) is decided inside build, --test, a start and a reload all run the same validation. Only the bind itself happens after build.
Inbounds that own their handle release it before the next generation binds. app/src/serve.rs → run_hysteria_inbound awaits Hy2Inbound::shutdown after its token is cancelled, because dropping a QUIC endpoint does not free its UDP port while connections are still live. Its comment calls this “what makes a reload transactional for an inbound that owns its handle”.
katana applies the same split in three places:
- Config reload.
src/runtime.rs→apply_reloadbuilds the new outbound pool when it changed, every node the reload adds, and the panel client of every running node whose entry changed, before it applies anything. If any of them fails to build, it logs “keeping current config” and returns: no node is removed, and no edit made alongside is applied. A running node whose entry changed then receives it as aStaticUpdate::Config, andsrc/manager/node.rs→NodeManager::apply_statictakes it “whole or not at all”: it builds the new panel client and router before it stores anything, and refuses an edit that does not build (“config edit refused, keeping the running one”). - Listener start.
src/manager/transport.rsbuilds the transport, the protocol tables or the Hysteria 2 server before it binds: “Everything that can fail happens before the bind, so a bad config never half-binds”. - User refresh.
src/manager/proxy.rs→ProxyManager::refreshbuilds the replacement user tables first. A build error is logged (“proxy refresh build failed, keeping current”) and “leaves the running state entirely untouched”. Only then does it publish the user registry and swap the tables.
The user registry itself is split into two phases in src/traffic.rs:
impl NodeTraffic { pub fn prepare(&self, entries: Vec<UserEntry>) -> PreparedUsers pub fn commit(&self, prepared: PreparedUsers)}prepare stages the next user set without changing the registry or reading any byte totals, so the next tables can be built, and can fail, before anything changes. commit swaps the set in one step and parks every outgoing counter as draining. It takes the registry lock and the draining lock together, in the order snapshot takes them, so a snapshot sees each outgoing counter in exactly one place, “never both, which would report its bytes twice and then commit them twice”.
| Test | File |
|---|---|
a_reload_rebinds_the_udp_port |
app/tests/integration/e2e_hysteria_inbound.rs |
a_reload_with_a_node_that_does_not_build_changes_nothing |
katana tests/unit/runtime.rs |
rate_change_drains_old_counter_and_reports_once |
katana tests/unit/traffic.rs |
rebound_credential_reports_the_old_uid_separately |
katana tests/unit/traffic.rs |
Deterministic time in tests
Section titled “Deterministic time in tests”Proxy timeouts are measured in seconds or minutes: a 10 s handshake limit, a 300 s idle limit, a one-second token bucket. A test that sleeps through them is slow and flaky. The code avoids real sleeps in three ways:
- Deadlines are events. A core test delivers
Event::Deadlineby hand, as inhandshake_deadline_fails_the_connectionabove. No clock is involved. - The runtime arms tokio timers.
ProxyServerRuntime::set_deadlinecomputestokio::time::Instant::now() + afterand usestokio::time::sleep_until, neverstd::time. Its comment gives the reason: “so a paused test clock (the only way to test a minutes-long idle timeout) is honoured”. A test marked#[tokio::test(start_paused = true)]then advances time instantly and exactly. - Wall clocks are injected. Timestamp checks take
now: fn() -> …in the constructor, as shown above.
katana’s TokenBucket also reads tokio::time::Instant, so its tests assert exact durations:
#[tokio::test(start_paused = true)]async fn a_chunk_larger_than_the_burst_is_still_limited() { let b = TokenBucket::new(10_000); let start = tokio::time::Instant::now(); b.consume(30_000).await; // 10 000 came out of the burst; the other 20 000 take two seconds. assert_eq!(start.elapsed(), Duration::from_secs(2));}| Test | File | Clock |
|---|---|---|
deadline_is_armed_against_the_tokio_clock |
concepts/tests/runtime.rs |
Paused; advances an hour first, so a deadline armed against std::time would fire at the wrong moment |
deadline_event_lets_the_core_time_out |
concepts/tests/runtime.rs |
Paused |
liveness_restarts_the_idle_deadline_on_progress |
protocols/tests/unit/transports/grpc_liveness.rs |
Paused |
handshake_deadline_fails_the_connection |
protocols/tests/unit/trojan/core.rs |
None; the deadline is an event |
an_idle_bucket_banks_one_second_and_no_more |
katana tests/unit/traffic.rs |
Paused |
the_limit_holds_however_the_writes_are_sized |
katana tests/unit/meter.rs |
Paused |
Panic-free parsing
Section titled “Panic-free parsing”Every byte a protocol parser reads comes from a peer that is not trusted. A panic in a parser kills the connection’s task, and a slice index computed from an attacker-controlled length is exactly where panics come from. So the library crates that parse wire data forbid the panicking forms at compile time. protocols/src/lib.rs and environment/src/lib.rs both carry, at the crate root:
#![deny( clippy::unwrap_used, clippy::expect_used, clippy::indexing_slicing, clippy::arithmetic_side_effects)]A #![cfg_attr(test, allow(…))] right after it lifts the rule for test code, which works on known-good inputs. etemenanki-concepts and etemenanki-app do not carry these crate-level denies. The validation gate maintainers run before a change lands is cargo clippy --workspace --all-targets --all-features -- -D warnings, so every other clippy warning fails the change as well.
protocols/src/helpers/parse.rs provides the replacements:
pub fn take<'a, I>(data: &'a [u8], index: I, what: &'static str) -> Result<&'a [u8], ProtocolError>where I: SliceIndex<[u8], Output = [u8]>,
pub fn take_array<const N: usize>(data: &[u8], at: usize) -> Result<[u8; N], ProtocolError>
pub fn need_more<T>(result: std::io::Result<T>) -> std::io::Result<Option<T>>takereturnsProtocolError::Truncated(what)instead of panicking on an out-of-range slice.take_arraychecks the offset addition (Overflow("field offset")) and the bounds, foru16::from_be_bytesand similar reads.need_moreseparates “the buffer is still short” from “the input is malformed”, so a parser that is handed a growing buffer can ask to be called again.
| Test | File |
|---|---|
varint_above_62_bits_is_refused_not_panicked |
protocols/tests/unit/hysteria/protocol.rs |
varint_truncated_at_every_boundary |
protocols/tests/unit/hysteria/protocol.rs |
tcp_response_truncated_at_every_boundary |
protocols/tests/unit/hysteria/protocol.rs |
udp_truncated_in_the_header_is_refused_and_never_panics |
protocols/tests/unit/hysteria/protocol.rs |
rejects_truncated_metadata |
protocols/tests/unit/mux/frame.rs |
Review rules
Section titled “Review rules”Maintainers check every change against the rules below. Each rule describes what the code must do; none of them is a statement about which parts of the current code already do it. A change that touches the area a rule covers should come with a test for the rule, including a negative test where the rule forbids something.
UDP associations accept only their own client
Section titled “UDP associations accept only their own client”A UDP association (a SOCKS UDP ASSOCIATE, or any relay that hands a client a UDP endpoint) accepts datagrams only from the client that set it up, and checks every packet against that client. On the client side, a link accepts replies only from the relay it talks to and never treats a datagram from another source as a relay reply. Otherwise any host that can reach the relay port can inject packets into someone else’s flow.
A change here needs a test in which a third-party address sends to the relay and the packet is dropped.
Per-connection queues are bounded, and pressure reaches the producer
Section titled “Per-connection queues are bounded, and pressure reaches the producer”Every queue that holds data for one connection has an explicit capacity. Moving items from a bounded channel into an unbounded queue does not count: it removes the bound the channel was there to provide. When a consumer stops, the pressure travels back to whoever produces the data, and memory does not grow. In the one-task runtime the fixed buffers are meant to be that bound; any component that adds a queue of its own needs a bound of its own.
The WireGuard driver (protocols/src/wireguard/device.rs) is a worked example. One driver task serves every connection in the tunnel, and each connection hands it application bytes over a bounded channel of CHANNEL_CAP (256) items. The driver takes at most one item ahead of a connection’s smoltcp socket and reads that connection’s channel only while it holds none. When the remote or the tunnel stops draining one connection, its socket fills, then its channel, and then its writer waits; the driver queues nothing more, and the other connections keep moving. On the way down, the driver reads a socket only while the connection’s channel toward the application has room.
| Test | File | What it pins |
|---|---|---|
a_held_item_keeps_the_channel_unread |
protocols/tests/unit/wireguard/device.rs |
While the driver holds an item, the channel is not read, so it fills and the producer’s next send fails with Full. |
a_stalled_tcp_flow_blocks_its_writer |
protocols/tests/pipeline/wireguard.rs |
A flow whose remote reads nothing blocks its writer after at most 512 KiB, another flow through the same tunnel still moves, and the waiting write goes through once the remote reads again. |
An active-connection limit covers the whole relay
Section titled “An active-connection limit covers the whole relay”A semaphore or permit that is described as limiting active connections is held from admission until the relay ends. A permit that covers only the handshake, the decode step or the transport step limits that stage, not active connections. At every spawn boundary, the permit has to move into the task that runs the relay and must not be dropped before the relay starts. A test for such a limit checks both sides: that the permit bounds something while the relay runs, and that it is released when the relay ends.
UDP to a multi-address name picks a usable family
Section titled “UDP to a multi-address name picks a usable family”When DNS returns several addresses for a UDP destination, the code picks one whose family (IPv4 or IPv6) matches a socket that is actually available. Taking the first answer and dropping the packet because that family has no socket ignores addresses that would have worked.
Fixed protocol fields are validated strictly
Section titled “Fixed protocol fields are validated strictly”Version bytes, reserved bytes, lengths, commands and type codes are checked against the values the protocol defines, and anything else is an error. The peer is not assumed to be a correct implementation. A change to a parser comes with negative tests that feed each fixed field a value the protocol does not define.
Configuration typos fail closed
Section titled “Configuration typos fail closed”Stable configuration structures reject unknown fields. This matters most for listen, port, protocol, network, security, TLS settings, outbounds and routing. An unknown value for a setting that selects behaviour is an error: a security value that is not recognised never results in plaintext because it is “not tls”.
Reloads are transactional
Section titled “Reloads are transactional”A reload does not destroy the running generation and only then discover that the new one cannot start. The intended order is:
- parse, validate and build the new configuration;
- create as many of the new resources as possible in advance;
- confirm that the new generation works;
- switch generations;
- on failure, keep or restore the old generation.
A failed reload also needs a way to be retried once an external dependency recovers. Waiting for the configuration file’s bytes to change is not enough.
katana accounting never loses or double-commits counters
Section titled “katana accounting never loses or double-commits counters”Traffic counters keep these properties:
- a failed report does not lose bytes;
- a retried report does not bill bytes twice;
- bytes of users who left (residuals) can be recovered and reported later;
- what happens when a user is removed and added back is defined explicitly;
- bytes added while a report is in flight are not committed by that report.
A committed report subtracts the amounts it reported rather than resetting a counter to zero, and a failed report hands its rows back so the next report retries them.
Rate limits charge actual bytes as debt
Section titled “Rate limits charge actual bytes as debt”A token bucket charges the number of bytes actually moved. A chunk larger than the bucket’s capacity is charged in full and the bucket goes into debt; the chunk does not pass after a single refill period. The effective rate does not depend on the size of the relay’s reads and writes.
Logs never contain credentials
Section titled “Logs never contain credentials”Log lines do not contain passwords, UUIDs, tokens, keys, panel secrets, or full URLs whose query string carries a secret. When a credential is malformed, the log names a safe identifier and the kind of error, never the value itself. Changes to HTTP error handling keep URLs and tokens redacted.