Skip to content

Design principles

Source files: 70 · checked against Etemenanki 596916d · katana v3.0.1
  • Etemenanki/concepts/src/lib.rs
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/wake.rs
  • Etemenanki/concepts/src/client.rs
  • Etemenanki/concepts/tests/runtime.rs
  • Etemenanki/concepts/tests/client.rs
  • Etemenanki/environment/src/lib.rs
  • Etemenanki/protocols/src/lib.rs
  • Etemenanki/protocols/src/error.rs
  • Etemenanki/protocols/src/helpers/parse.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/core/harness.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/src/vmess/accounts.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/protocols/src/tun/config.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/src/wireguard/device.rs
  • Etemenanki/protocols/src/transports/grpc/stream.rs
  • Etemenanki/protocols/src/hysteria/connection.rs
  • Etemenanki/protocols/src/hysteria/server/config.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/datagrams.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/tests/unit/trojan/core.rs
  • Etemenanki/protocols/tests/unit/vmess/core.rs
  • Etemenanki/protocols/tests/unit/vmess/protocol.rs
  • Etemenanki/protocols/tests/unit/ss_2022/core.rs
  • Etemenanki/protocols/tests/unit/mux/demux.rs
  • Etemenanki/protocols/tests/unit/wireguard/device.rs
  • Etemenanki/protocols/tests/pipeline/wireguard.rs
  • Etemenanki/protocols/tests/unit/mux/frame.rs
  • Etemenanki/protocols/tests/unit/hysteria/protocol.rs
  • Etemenanki/protocols/tests/unit/hysteria/server/datagrams.rs
  • Etemenanki/protocols/tests/unit/transports/grpc_liveness.rs
  • Etemenanki/app/src/config.rs
  • Etemenanki/app/src/transport.rs
  • Etemenanki/app/src/inbound/mod.rs
  • Etemenanki/app/src/inbound/tun.rs
  • Etemenanki/app/src/outbound/mod.rs
  • Etemenanki/app/src/main.rs
  • Etemenanki/app/src/instance.rs
  • Etemenanki/app/src/serve.rs
  • Etemenanki/app/src/outbound/udp_fanout.rs
  • Etemenanki/app/tests/unit/config.rs
  • Etemenanki/app/tests/unit/transport.rs
  • Etemenanki/app/tests/unit/inbound.rs
  • Etemenanki/app/tests/integration/e2e_hysteria.rs
  • Etemenanki/app/tests/integration/e2e_hysteria_inbound.rs
  • katana/src/serve.rs
  • katana/src/connector.rs
  • katana/src/meter.rs
  • katana/src/traffic.rs
  • katana/src/config.rs
  • katana/src/manager/proxy.rs
  • katana/src/manager/transport.rs
  • katana/src/manager/node.rs
  • katana/src/runtime.rs
  • katana/tests/unit/meter.rs
  • katana/tests/unit/traffic.rs
  • katana/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.

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

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.

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.

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:

protocols/src/vmess/core.rs
impl<T> VMessCore<T> {
pub fn new(
validator: Arc<AccountValidator<T>>,
now: fn() -> i64,
sniff: bool,
source: Option<IpAddr>,
) -> Self
}
protocols/src/ss_2022/core.rs
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.

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.

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:

protocols/tests/unit/trojan/core.rs
#[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 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.

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 + 1 bytes free. A datagram outbound is read only while staging has STAGING_RESERVE plus 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.

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 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.

Every per-connection resource has a size or a count that is fixed in the code and visible in review.

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 example drive::<{ 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.

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.

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

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”.

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.toml
ERROR 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.”

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

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_reload builds 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 a StaticUpdate::Config, and src/manager/node.rs → NodeManager::apply_static takes 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.rs builds 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::refresh builds 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

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:

  1. Deadlines are events. A core test delivers Event::Deadline by hand, as in handshake_deadline_fails_the_connection above. No clock is involved.
  2. The runtime arms tokio timers. ProxyServerRuntime::set_deadline computes tokio::time::Instant::now() + after and uses tokio::time::sleep_until, never std::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.
  3. 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:

katana tests/unit/traffic.rs
#[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

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>>
  • take returns ProtocolError::Truncated(what) instead of panicking on an out-of-range slice.
  • take_array checks the offset addition (Overflow("field offset")) and the bounds, for u16::from_be_bytes and similar reads.
  • need_more separates “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

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.

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”.

A reload does not destroy the running generation and only then discover that the new one cannot start. The intended order is:

  1. parse, validate and build the new configuration;
  2. create as many of the new resources as possible in advance;
  3. confirm that the new generation works;
  4. switch generations;
  5. 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.

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.

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.