Skip to content

Server cores

Source files: 33 · checked against Etemenanki 596916d
  • Etemenanki/concepts/src/core.rs
  • Etemenanki/concepts/src/buffer.rs
  • Etemenanki/concepts/src/runtime.rs
  • Etemenanki/concepts/tests/runtime.rs
  • Etemenanki/protocols/src/core/mod.rs
  • Etemenanki/protocols/src/core/harness.rs
  • Etemenanki/protocols/tests/unit/core/mod.rs
  • Etemenanki/protocols/src/flow.rs
  • Etemenanki/protocols/src/mux/demux.rs
  • Etemenanki/protocols/tests/unit/mux/demux.rs
  • Etemenanki/protocols/src/ss_legacy/aead.rs
  • Etemenanki/protocols/src/vmess/framing.rs
  • Etemenanki/protocols/src/sniff/mod.rs
  • Etemenanki/protocols/src/sniff/collector.rs
  • Etemenanki/protocols/src/http/core.rs
  • Etemenanki/protocols/src/http/protocol.rs
  • Etemenanki/protocols/src/ss_legacy/core.rs
  • Etemenanki/protocols/tests/unit/ss_legacy/aead.rs
  • Etemenanki/protocols/src/ss_2022/core.rs
  • Etemenanki/protocols/src/trojan/core.rs
  • Etemenanki/protocols/src/trojan/protocol.rs
  • Etemenanki/protocols/src/vless/core.rs
  • Etemenanki/protocols/src/vmess/core.rs
  • Etemenanki/protocols/tests/unit/vmess/core.rs
  • Etemenanki/protocols/src/hysteria/protocol.rs
  • Etemenanki/protocols/src/hysteria/server/inbound.rs
  • Etemenanki/protocols/src/hysteria/server/datagrams.rs
  • Etemenanki/protocols/src/tun/udp.rs
  • Etemenanki/protocols/src/tun/inbound.rs
  • Etemenanki/protocols/src/socks/server.rs
  • Etemenanki/protocols/src/socks/protocol.rs
  • Etemenanki/protocols/tests/unit/socks/server.rs
  • Etemenanki/app/src/serve.rs

Every server-side protocol in Etemenanki except SOCKS is a core: a plain Rust value that implements ProxyCoreDecode and never touches a socket, a clock or a task Context. The per-connection runtime owns the I/O. It reads bytes, hands them to the core one Event at a time, and carries out the Effects the core pushes back. Because a core is a function from bytes and events to bytes and effects, a test can drive it with hand-built input and assert on the output. SOCKS is the exception because its UDP side lives on a second socket, a relay “hub” that the control connection only keeps alive; protocols/src/socks/server.rs runs it as one task outside this contract. That task also decides whom the hub hears: ExpectedSender holds each association to the control connection’s client, requiring the peer’s IP (compared in canonical form by endpoint in protocols/src/socks/protocol.rs) on every datagram and pinning the port of the first datagram it forwards, or a non-zero port the UDP ASSOCIATE names alongside that same IP. Over a Unix socket, which has no peer, the request must name the exact source address and port. Datagrams from any other address are dropped unread. only_the_control_peer_is_heard_and_its_first_datagram_pins_the_port in protocols/tests/unit/socks/server.rs pins the rule, and SOCKS has the full table.

This page is the contract between a core and the runtime: the trait, every event and effect, the Effects sink, and the rules that make in-place decryption, held buffers, keys, deadlines and datagrams safe. Read it before you write a new core or change the delivery logic in concepts/src/runtime.rs. The runtime’s scheduling, wakers and backpressure are covered in Server runtime; the client side (ProxyCoreEncode) is covered in Client runtime.

A core decides; the runtime acts.

The core The runtime (ProxyServerRuntime)
Parses the client’s wire bytes, decrypting in place Reads the transport into a fixed read buffer
Decides where each flow goes and mints a key for it (Effect::Open) Dials the target through its Connector and tracks one slot per key
Names which byte ranges go to which outbound (Effect::Forward) Writes those ranges, in order, with backpressure
Seals reply bytes into the staging area Writes the staging area to the transport
Keeps its own timers and arms the earliest one (Effect::SetDeadline) Owns the single tokio timer and delivers Event::Deadline
Says when the connection is over (Effect::Finish) Drains, shuts down and resolves to the Traffic totals

Randomness, wall-clock time and shared account state are constructor arguments of the concrete core, never something the runtime injects. That keeps every call deterministic for a given core value. VMessCore::new takes now: fn() -> i64, Ss2022Core::new takes now: fn() -> u64 (with Ss2022Core::with_system_clock as the production constructor), and Hy2UdpCore counts its own sweep ticks as its clock.

All of the contract lives in concepts/src/core.rs. The shared building blocks that the in-tree cores compose live in protocols/src/core/mod.rs.

concepts/src/core.rs
pub struct Single;
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] {
&[]
}
}
Item Meaning
Key Names one outbound. The core mints keys with Effect::Open, and every outbound-side event carries the key it came from. A core with exactly one outbound uses the unit struct Single. The runtime stores slots in a BTreeMap, hence Ord.
Target What the runtime’s Connector dials. For the in-tree server cores it is Flow<T> (destination, user, sniff result, source address).
Error Returned from handle; the runtime wraps it as RuntimeError::Core and ends the connection. The in-tree cores use io::Error.
TransportAddr The peer address of a datagram transport packet: what Event::TransportDatagram arrives from and Effects::stage_to sends to. SocketAddr for one UDP socket serving many peers, () for a QUIC connection’s datagrams (one peer) and for any byte-stream core.
STAGING_RESERVE Upper bound on bytes the core stages toward the transport while handling one event, beyond echoing that event’s own payload. See The staging reserve.
MAX_DATAGRAM The largest datagram delivered whole, from a datagram outbound or a datagram transport. Longer packets are truncated the way a kernel recv truncates. Default 4096.
handle Advance the state machine by one event. Returns how many leading bytes of a byte-carrying event were consumed.
held Bytes the core keeps on its own account and forwards by range with ForwardHeld or SendToHeld. Empty by default.
concepts/src/core.rs
pub enum Event<'a, C: ProxyCoreDecode + ?Sized> {
Transport(&'a mut [u8]),
TransportDatagram {
from: C::TransportAddr,
data: &'a mut [u8],
},
TransportSendFailed {
to: C::TransportAddr,
error: io::Error,
},
TransportEof,
Outbound { key: C::Key, data: &'a mut [u8] },
Datagram {
key: C::Key,
from: Destination,
data: &'a mut [u8],
},
SendFailed {
key: C::Key,
to: Destination,
error: io::Error,
},
OutboundEof { key: C::Key },
Connected { key: C::Key },
ConnectFailed { key: C::Key, error: io::Error },
OutboundError { key: C::Key, error: io::Error },
Deadline,
}

The Debug impl prints slice lengths, never slice contents.

Event When the runtime delivers it What handle must return
Transport(data) A byte-stream transport produced bytes. data is the whole unparsed region of the read buffer, including any tail the core left on the previous call. Bytes consumed from the front, 0..=data.len(). The rest is offered again.
TransportDatagram { from, data } A datagram transport produced one whole packet. Never delivered over a stream transport. Exactly data.len().
TransportSendFailed { to, error } A packet staged toward peer to was refused by the link and dropped. Not fatal. 0 by convention; the count is not checked.
TransportEof The stream transport’s read side ended. A datagram transport never delivers it. 0 by convention.
Outbound { key, data } Stream outbound key produced bytes (read into the scratch buffer). Exactly data.len(): there is no per-outbound buffer to keep a remainder in.
Datagram { key, from, data } Datagram outbound key produced one packet. from is the packet’s source; a link that resolves domains reports the address it received from, as an IP. Exactly data.len().
SendFailed { key, to, error } An Effect::SendTo or SendToHeld on datagram outbound key was refused. The packet is dropped and the key stays live. 0 by convention.
OutboundEof { key } Outbound key’s read side ended. 0 by convention.
Connected { key } The dial started by Effect::Open for key completed. 0 by convention.
ConnectFailed { key, error } The dial failed. The key is already gone. 0 by convention.
OutboundError { key, error } Outbound key failed while being read, written, flushed or shut down. The key is already gone. 0 by convention.
Deadline The deadline armed by Effect::SetDeadline passed. It is now disarmed. 0 by convention.

The runtime enforces the consume rules in feed_transport, poll_transport_recv_packet and after_scratch: consuming more than the slice, or less than a whole outbound payload or transport datagram, ends the connection with RuntimeError::BadConsume.

concepts/src/core.rs
pub enum Effect<C: ProxyCoreDecode + ?Sized> {
Forward { key: C::Key, range: Range<usize> },
SendTo {
key: C::Key,
to: Destination,
range: Range<usize>,
},
ForwardHeld { key: C::Key, range: Range<usize> },
SendToHeld {
key: C::Key,
to: Destination,
range: Range<usize>,
},
Open { key: C::Key, target: C::Target },
Shutdown { key: C::Key },
Close { key: C::Key },
ShutdownTransport,
SetDeadline(Option<Duration>),
Finish,
}

Effects are applied strictly in the order they were pushed, by ProxyServerRuntime::drive_effects. An effect that cannot complete yet stays at the front of the queue, and everything behind it waits, whichever key it names. The runtime keeps the event’s bytes alive meanwhile and reads nothing new into that buffer.

Effect What applying it does Can it block the queue?
Forward Writes range of the event’s slice to stream outbound key, then flushes it. A short write advances range.start and writes again at once, until the write returns Pending. A write of zero bytes fails the outbound with WriteZero. An empty range is dropped. Yes: while key is still connecting, or while the write returns Pending.
SendTo Sends range of the event’s slice as one datagram to to through datagram outbound key. Yes: while connecting, or while the send returns Pending.
ForwardHeld As Forward, but range indexes held(), resolved at the moment of writing. Yes, as Forward.
SendToHeld As SendTo, but range indexes held(). Yes, as SendTo.
Open Calls connector.connect(target), stores the dial future in a new slot for key and schedules the key. No. The runtime polls the dial future when the key is scheduled, alongside its other work; forwards to key behind it wait for the connect. An Open queued behind a blocked effect is not applied, and its dial not started, until that effect completes.
Shutdown Half-closes the write side of stream outbound key once everything queued before it is written. On a datagram outbound it completes at once. Yes: while connecting, or while poll_shutdown returns Pending.
Close Drops the slot for key in both directions, including a dial still in flight. No event follows. No.
ShutdownTransport Sets a flag. The runtime half-closes the transport’s write side after staging has drained and been flushed. No.
SetDeadline(after) Some(d) arms the single timer for d from now, replacing any earlier deadline. None disarms it. No.
Finish Marks the connection over. The runtime stops reading, writes out staging, honours a pending ShutdownTransport, and resolves. No.
concepts/src/core.rs
pub const INLINE_EFFECTS: usize = 4;
pub type EffectList<C> = SmallVec<[Effect<C>; INLINE_EFFECTS]>;
pub type PacketList<A> = VecDeque<(usize, A)>;
pub struct Effects<'a, C: ProxyCoreDecode + ?Sized> { /* private */ }
impl<'a, C: ProxyCoreDecode + ?Sized> Effects<'a, C> {
pub fn new(list: &'a mut EffectList<C>, staging: Staging<'a>) -> Self;
pub fn with_packets(
list: &'a mut EffectList<C>,
staging: Staging<'a>,
packets: &'a mut PacketList<C::TransportAddr>,
) -> Self;
pub fn push(&mut self, effect: Effect<C>);
pub fn staging(&mut self) -> &mut Staging<'a>;
pub fn stage(&mut self, bytes: &[u8]) -> Option<()>;
pub fn transport_is_datagram(&self) -> bool;
pub fn stage_to(&mut self, to: C::TransportAddr, len: usize) -> Option<&mut [u8]>;
pub fn put_to(&mut self, to: C::TransportAddr, bytes: &[u8]) -> Option<()>;
pub fn list(&self) -> &[Effect<C>];
}
Method Use
push Append an effect. The list is owned by the runtime, reused across calls, and inline up to INLINE_EFFECTS (a handshake is typically Open + Forward + SetDeadline).
staging The transport-bound Staging area. Staging::reserve(len) claims len bytes at the tail and returns them for filling, which is how a frame is sealed in place with no intermediate Vec. Staging::room says how much is free.
stage Copy bytes toward a byte-stream transport. None when the room is too small.
transport_is_datagram true when the sink was built with with_packets, so staged bytes need a peer.
stage_to Claim len bytes as one packet toward to and return them for filling. None if there is not enough room, or if the transport is a byte stream.
put_to Copy bytes as one packet toward to.
list The effects pushed so far, for a core that inspects its own output in tests.

Staging is the append-only view of the runtime’s transport-bound WriteBuffer<BUF_SIZE> that the sink wraps:

concepts/src/buffer.rs
pub struct Staging<'a> { /* private */ }
impl Staging<'_> {
pub fn room(&self) -> usize;
pub fn reserve(&mut self, len: usize) -> Option<&mut [u8]>;
pub fn put(&mut self, bytes: &[u8]) -> Option<()>;
}

reserve claims len bytes at the tail and counts them as staged at once; the core fills them in place. put is reserve plus a copy, and Effects::stage is put. Staging never partially commits: a reserve that does not fit returns None and stages nothing (staging_refuses_over_reservation_without_partial_commit in concepts/src/buffer.rs). Before it hands out a Staging, WriteBuffer::staging calls WriteBuffer::room, which slides still-pending bytes to the front once the free tail is used up; that slide is the only copy the staging buffer makes on its own (write_buffer_slides_pending_when_full).

Build with Effects::new. Append with stage or staging().reserve(..). The runtime writes the whole pending region of its WriteBuffer to the transport.

One connection is one runtime and one core. The runtime calls handle synchronously from inside its poll, so a core never runs concurrently with itself and needs no locks.

sequenceDiagram
  participant T as Transport
  participant R as ProxyServerRuntime
  participant C as Core
  participant O as Outbound
  T->>R: bytes into the read buffer
  R->>C: handle(Event::Transport(unparsed))
  C-->>R: consumed, Open(k) + Forward(k, range)
  R->>R: validate ranges, queue in order
  R->>O: Open applied: connector.connect(target)
  Note over R: Forward waits while k is connecting
  O-->>R: dial completes
  R->>C: handle(Event::Connected(k))
  R->>O: Forward applied: write range
  O-->>R: reply bytes into scratch
  R->>C: handle(Event::Outbound(k, data))
  C-->>R: sealed reply in staging
  R->>T: write staging

Every delivery goes through the same three steps, whatever the event. ProxyServerRuntime::deliver builds the sink and calls handle; enqueue validates each ranged effect, rebases event-slice ranges onto the buffer they index and appends the effect to the queue; drive_effects then applies the queue from the front until an effect blocks. Effects pushed while answering a notification (Connected, Deadline, OutboundError and the rest) join the back of the same queue.

flowchart TB
  ev["Event ready and allowed"] --> deliver["deliver: core.handle(event, effects)"]
  deliver --> check{"consume count valid?"}
  check -- no --> bad["RuntimeError::BadConsume"]
  check -- yes --> enq["enqueue: check ranges, rebase, tag source buffer"]
  enq --> oob{"range in bounds?"}
  oob -- no --> rng["RuntimeError::RangeOutOfBounds"]
  oob -- yes --> drive["drive_effects: apply the front effect"]
  drive --> blocked{"front effect blocked?"}
  blocked -- "no, more queued" --> drive
  blocked -- "yes, or queue empty" --> wait["return to the scheduler"]

concepts/tests/runtime.rs drives a toy multiplexing protocol, TinyMux, whose frames are [kind][key][len:u16][payload] with an XOR “cipher” (XOR = 0x55, HDR = 4). It declares Key = u8, Target = String, TransportAddr = () and STAGING_RESERVE = HDR + 64. Its transport arm is the canonical parse loop:

concepts/tests/runtime.rs (excerpt)
Event::Transport(buf) => {
let mut consumed = 0;
loop {
let rest = &mut buf[consumed..];
if rest.len() < HDR {
break;
}
let (k, key) = (rest[0], rest[1]);
let len = u16::from_be_bytes([rest[2], rest[3]]) as usize;
if rest.len() < HDR + len {
break;
}
let payload = &mut rest[HDR..HDR + len];
for b in payload.iter_mut() {
*b ^= XOR;
}
let (start, end) = (consumed + HDR, consumed + HDR + len);
match k {
kind::OPEN => { /* insert key, push Effect::Open */ }
kind::DATA => fx.push(Effect::Forward { key, range: start..end }),
kind::CLOSE => fx.push(Effect::Shutdown { key }),
kind::FIN => { /* ShutdownTransport, Finish */ }
other => return Err(format!("bad frame kind {other}")),
}
consumed = end;
}
/* optional idle SetDeadline */
Ok(consumed)
}

core_is_driven_without_any_io drives it with no runtime at all:

  1. Input. The test builds OPEN 9 "somewhere", DATA 9 "hi" and the first 5 bytes of DATA 9 "split", and calls handle(Event::Transport(&mut input), &mut fx) over an Effects::new sink on a WriteBuffer::<256>.

  2. Parse frame by frame. The core decodes the two complete frames and stops at the incomplete third one. It returns HDR + 9 + HDR + 2, the offset where the split frame starts. The runtime would keep those 5 bytes and offer them again, prefixed to the next read.

  3. Effects. The list holds Open { key: 9, target: "somewhere" } then Forward { key: 9, range: 17..19 }. The range is relative to the slice the core was given; the runtime rebases it onto the read buffer when it queues the effect.

  4. In place. input[17..19] now reads hi: the payload was decrypted where it lay, and the forward will write it with no copy.

  5. Connected. handle(Event::Connected { key: 9 }) stages the frame [5, 9, 0, 0] (kind::CONNECTED). No effect is needed.

  6. Downlink. handle(Event::Outbound { key: 9, data: b"pong" }) seals a DATA frame into staging with fx.staging().reserve(HDR + len) and returns 4, the whole payload.

The same core then runs under a real ProxyServerRuntime over tokio::io::duplex pipes: relays_two_keys_and_completes_on_fin opens two keys and checks the Traffic totals, and stalled_outbound_holds_uplink_but_not_other_downlink gives key 1 an 8-byte pipe and checks that key 2’s downlink still arrives while the uplink behind key 1’s forward waits.

Event::Transport hands the core the whole unparsed region [start..end) of the ReadBuffer. The core returns how much of its front it consumed. The runtime (feed_transport) advances start by that count and calls again while the core makes progress, bytes remain, and no queued effect pins the read buffer or the held buffer. It reads more from the transport once the unparsed region is empty or a call consumed nothing. Parse frame by frame, stop at the first incomplete frame, and return the offset of its start.

  • Mechanism: ReadBuffer keeps parsed bytes [0..start) in place, since a queued forward may still reference them. The runtime reads nothing into the buffer while a queued forward pins it, and calls compact to move the unparsed tail to the front only when the free tail is used up.
  • Tests: core_is_driven_without_any_io (the split frame stays unconsumed) in concepts/tests/runtime.rs; read_buffer_compacts_unparsed_tail_to_front in concepts/src/buffer.rs.

The slice is &mut, so a frame can be decrypted in place and forwarded by range. The consequence: if a frame’s header is decrypted in place but its body has not arrived, the same bytes come back on the next call already transformed, and a naive parser decrypts them twice. Keep the decoded length, and any cipher state that advanced (a nonce counter, a SHAKE mask), in the core, and skip the header on the next call.

  • Mechanism: ChunkDecoder in protocols/src/ss_legacy/aead.rs opens the length from a copy, increments its NonceCounter, and stores pending_len: Option<usize>. The payload is opened in place only once the whole chunk is present. protocols/src/vmess/framing.rs has its own ChunkDecoder for VMess.
  • Tests: a_chunk_split_across_reads_decodes_its_header_once in protocols/tests/unit/vmess/core.rs; chunk_encoder_and_decoder_agree_and_never_decrypt_twice in protocols/tests/unit/ss_legacy/aead.rs (the same header is offered three times, with more of the body each time).

BUF_SIZE is a const generic of ProxyServerRuntime and sizes each of its three buffers: transport read, transport staging and outbound scratch. When the unparsed region fills the whole read buffer and the core still consumes nothing, the runtime fails with RuntimeError::FrameTooLarge. Size BUF_SIZE to the largest frame the protocol admits: an AEAD chunk plus its overhead, a full HTTP head, a mux frame with its metadata. Each in-tree core publishes the value its runtime needs as an associated BUF_SIZE constant (see Limits).

  • Mechanism: when the read buffer has no free tail left, poll_transport_read checks ReadBuffer::is_saturated (start == 0 && end == N) before it calls compact. A saturated buffer is one that compacting cannot free, so the runtime returns FrameTooLarge.
  • Tests: frame_larger_than_the_buffer_is_an_error (a 500-byte frame into BUF_SIZE = 128) in concepts/tests/runtime.rs; read_buffer_saturation_means_frame_too_large in concepts/src/buffer.rs.

ProxyServerRuntime::build also asserts BUF_SIZE > Core::STAGING_RESERVE and panics with “BUF_SIZE must exceed the core’s STAGING_RESERVE or nothing is ever read”.

A core may forward bytes that are not in the event’s slice: bytes it keeps itself and exposes through held(). Uses in the tree:

  • Sniffing. Copy a flow’s first bytes into the held buffer and consume them, arm SetDeadline, and keep collecting until a TLS SNI or HTTP Host is found, the budget is spent, or the deadline fires. Then push Open with the sniffed target followed by a ForwardHeld of everything collected. SniffPrefix does this for every in-tree core that sniffs a connection’s own flow. mux.cool sub-flows are the exception: Demux sniffs only the payload of the New frame and opens at once, because holding the carrier for one sub-flow would stall the others.
  • A rewritten request, such as a plain-HTTP proxy’s forwarded head, that must reach the outbound but never existed on the wire.
  • A frame reassembled across chunks, such as a mux frame that spans VMess chunks (Demux::feed_chunks), or a Hysteria UDP packet reassembled from fragments (Hy2UdpCore sends it with SendToHeld). Demux::feed_chunks takes every chunk one byte event opened in a single call, because each held forward it pushes reads the held buffer after the event, so the frames completed earlier in the event must stay in place until then.

The rule that makes this safe: a held range is resolved when its effect is applied, not when it is pushed. That may be several calls later if the outbound is still connecting or not writable. Until every queued held effect has been applied, the runtime delivers no byte event, so at the start of any byte event the core may clear or overwrite the held buffer freely. The notifications that can still arrive meanwhile may only append to it; a Vec that reallocates on append is fine, because the range is re-resolved against the current slice on every write attempt.

Event Raised from Staging room checked first Withheld while a queued effect pins…
Transport poll_transport_read, feed_transport STAGING_RESERVE, once at the top of each read pass the read buffer or the held buffer
TransportDatagram poll_transport_recv_packet STAGING_RESERVE the read buffer or the held buffer
TransportEof poll_transport_read STAGING_RESERVE the read buffer or the held buffer
Outbound, OutboundEof service_key STAGING_RESERVE + 1 the scratch buffer or the held buffer
Datagram service_key STAGING_RESERVE + datagram_limit() the scratch buffer or the held buffer
Connected, ConnectFailed service_key, while the key is connecting STAGING_RESERVE + 1 never withheld
OutboundError service_key (read or flush error) or drive_effects (write or shutdown error) only for a read error, which needs the read’s own room check to pass first a read error only surfaces from a read, which is withheld like Outbound; flush, write and shutdown errors are never withheld
SendFailed drive_effects not checked never withheld
TransportSendFailed poll_transport_send_packets not checked never withheld
Deadline poll_deadline not checked never withheld

The held buffer is the core’s to size, and it is the one per-connection allocation the runtime does not cap. Bound it by what the protocol admits; SniffPrefix is bounded by SNIFF_LIMIT.

  • Mechanism: Queued::pins(Source::Held) and ProxyServerRuntime::held_free, checked in feed_transport, poll_transport_read and service_key. Held ranges are validated against held().len() when queued (enqueue) and again when applied (drive_effects); out of range is RuntimeError::RangeOutOfBounds.
  • Tests: held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites (a gated dial keeps the ForwardHeld queued while more client bytes arrive; they reach the outbound only after it) and a_held_range_past_the_buffer_is_rejected, both in concepts/tests/runtime.rs; vmess_keeps_every_frame_one_read_completes in protocols/tests/unit/mux/demux.rs (one read opens three VMess chunks, two frames straddle chunk boundaries, and both are forwarded from the held buffer).

STAGING_RESERVE is the most a core stages in one call beyond the payload it is given: a reply header, or one frame’s overhead times the number of frames one outbound read may be split into. The runtime polls a stream outbound only when staging has STAGING_RESERVE + 1 bytes free, and reads at most room - STAGING_RESERVE bytes (capped at BUF_SIZE) into scratch. An Outbound event of n bytes therefore always comes with at least STAGING_RESERVE + n bytes of room, so sealing it never fails and Staging::reserve can be unwrapped, or mapped to the staging_full() error as Passthrough::on_outbound does.

A transport event only has the reserve checked, once, before the read; feed_transport then calls the core again on the unparsed tail as long as it makes progress, without checking again. Transport bytes are expected to be forwarded, not staged. A core that echoes transport bytes back, such as a server that answers a probe, reads Staging::room and consumes fewer bytes if it cannot stage its answer.

Deadline is delivered regardless of room, so that an idle timeout still fires against a peer that has stopped reading. The same holds for the notifications raised while effects are applied or packets sent (see the table above). A core must tolerate stage failing in those handlers.

  • Mechanism: the room_short check in service_key and the staging.room() < Core::STAGING_RESERVE check in poll_transport_read.
  • Tests: datagram_outbound_is_never_truncated_by_staging_backpressure in concepts/tests/runtime.rs; write_buffer_slides_pending_when_full in concepts/src/buffer.rs.

A key is live from Effect::Open until its outbound is gone. The runtime removes the slot on any of these:

stateDiagram-v2
  [*] --> Connecting: Effect Open applied
  Connecting --> Live: dial ok, Event Connected
  Connecting --> [*]: dial failed, Event ConnectFailed
  Live --> WriteClosed: Effect Shutdown applied
  Live --> ReadClosed: Event OutboundEof
  WriteClosed --> [*]: Event OutboundEof
  ReadClosed --> [*]: Effect Shutdown applied
  Live --> [*]: Effect Close, or Event OutboundError

Close also works while the key is connecting, and drops the dial. An I/O error in any state after the connect removes the key and delivers OutboundError.

  • Open on a live key fails the connection with RuntimeError::DuplicateKey (checked in drive_effects when the Open is applied).
  • An effect for an unknown key. Forward, SendTo, their held variants and Shutdown aimed at a key with no slot fail with RuntimeError::UnknownKey. A Forward or ForwardHeld with an empty range is dropped before the key is looked up, and Close of an unknown key is a no-op.
  • The wrong link kind. A stream effect on a datagram outbound, or the reverse, fails with RuntimeError::WrongLinkKind.
  • Once a key is gone after a failure, forget_key drops every queued effect still aimed at it (a forward waiting on a dial that just failed has nowhere to go), and the core is told once, through ConnectFailed or OutboundError.
  • A datagram outbound that refuses a send keeps its key; the core hears SendFailed. Only a failed receive removes a datagram key, since that is the socket or connection itself being gone.
  • A datagram outbound never reports end of stream. Shutdown on it completes at once and marks the write side closed, but no OutboundEof follows, so a datagram key lives until Close, a failed receive, or the end of the connection. Hy2UdpCore retires quiet sessions with Close from its sweep deadline; TunUdpCore, which serves one association per runtime, ends the whole runtime with Finish instead.
  • A key removed by the core itself (Close, or the second half-close) is gone with no event. Queued effects that still name it are not dropped, so pushing a Forward or Shutdown for it afterwards ends the connection with UnknownKey.

A protocol whose peer may reuse a session id right after closing it must either Close the old key before reopening, or tag keys with a generation. The mux.cool demultiplexer does the second. Its key is FlowKey::Sub(SubKey):

protocols/src/core/mod.rs
pub enum FlowKey {
Direct,
Sub(SubKey),
}
pub struct SubKey {
pub id: u16,
pub generation: u32,
}

Demux increments one generation counter per carrier (wrapping_add(1)) for every New frame it accepts. When the peer ends session 1 and immediately opens a new session 1, the new outbound gets a fresh SubKey while the old one may still be draining its read side, and downlink bytes that still arrive for the retired generation are framed to nobody.

  • Tests: stream_sessions_open_forward_and_end_with_fresh_generations in protocols/tests/unit/mux/demux.rs; connect_failure_reaches_the_core_as_an_event and a_refused_datagram_send_keeps_the_key_alive in concepts/tests/runtime.rs. No test drives DuplicateKey, UnknownKey, WrongLinkKind or BadConsume.

There is exactly one timer per connection. A protocol with several timeouts (a handshake limit, an idle limit per session, a keep-alive) keeps its own ordered map of expiries, arms the earliest with SetDeadline, and re-arms on every Deadline. A Deadline event means the timer fired and is now disarmed; arming again replaces any earlier deadline.

The runtime arms against tokio::time::Instant, not std::time::Instant, so a paused tokio clock in tests is honoured. The core must never call Instant::now() inside handle; wall-clock time comes from a clock passed to its constructor, so tests can drive it.

  • Mechanism: ProxyServerRuntime::set_deadline (tokio::time::sleep_until, or Sleep::reset on an existing timer) and poll_deadline, which step polls ahead of any effect or I/O work, so a connection that is always ready for I/O cannot starve its own deadline.
  • Tests: deadline_event_lets_the_core_time_out and deadline_is_armed_against_the_tokio_clock (the clock is advanced an hour before the runtime starts; a 5 s idle deadline must not fire within 4 s) in concepts/tests/runtime.rs.

Datagrams follow the socket rule on both sides: one packet per read into an empty buffer, consumed whole, and a packet longer than the limit truncated the way a kernel recv truncates.

Datagram outbound Datagram transport
Arrives as Event::Datagram { key, from, data } Event::TransportDatagram { from, data }
Sent with Effect::SendTo / SendToHeld Effects::stage_to / put_to
Read limit datagram_limit() = min(MAX_DATAGRAM, BUF_SIZE - STAGING_RESERVE) min(MAX_DATAGRAM, BUF_SIZE)
Read only when staging has STAGING_RESERVE + datagram_limit() free STAGING_RESERVE free
A refused send packet dropped, key kept, Event::SendFailed packet dropped, runtime carries on, Event::TransportSendFailed
Addressed by Destination, which may name a domain; the DatagramLink decides how to reach it TransportAddr, whatever the transport link addresses by

Because a datagram outbound is read only with room for a whole packet plus the reserve, a client that stops reading stops the downlink at the outbound socket instead of losing packet tails. Raise MAX_DATAGRAM for a protocol whose packets exceed an MTU, and size BUF_SIZE to match.

A datagram transport serves many peers on one runtime, so its core is a demultiplexer: derive the Key from the packet’s peer, Open on first sight, and expire idle peers from the deadline. There is no end of stream on that side; Effect::Finish is the only way such a runtime ends. A packet the transport refuses is dropped rather than failing the connection, because one peer’s unreachable address must not take the others down.

  • Tests: datagrams_round_trip_through_a_real_udp_socket, datagram_transport_demultiplexes_peers_and_frames_replies, staging_without_a_peer_under_a_datagram_transport_is_an_error, stage_to_is_refused_under_a_stream_transport, a_refused_transport_packet_is_reported_not_fatal and single_peer_datagram_transport_demuxes_sessions_and_reports_refusals, all in concepts/tests/runtime.rs.

The in-tree cores do not inherit from a base type. They compose three small helpers from protocols/src/core/mod.rs, are generic over the user payload T (the app uses ()), fail with io::Error, and expose is_established() so the application can tell a client that never spoke from one that is being served.

protocols/src/core/mod.rs
pub enum Phase { Handshake, Sniff, Relay, Closing }
pub enum Expired { Handshake, Sniff, Idle }
pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10);
pub const RELAY_IDLE_TIMEOUT: Duration = Duration::from_secs(300);
impl Timing {
pub fn new() -> Self;
pub fn phase(&self) -> Phase;
pub fn is_established(&self) -> bool;
pub fn touch<C: ProxyCoreDecode>(&mut self, fx: &mut Effects<'_, C>);
pub fn enter<C: ProxyCoreDecode>(&mut self, phase: Phase, fx: &mut Effects<'_, C>);
pub fn expired<C: ProxyCoreDecode>(&mut self, fx: &mut Effects<'_, C>) -> Expired;
}
stateDiagram-v2
  [*] --> Handshake
  Handshake --> Sniff: enter Sniff
  Handshake --> Relay: enter Relay
  Sniff --> Relay: enter Relay
  Relay --> Closing: idle deadline, Finish pushed
Phase Deadline armed On Deadline, expired returns
Handshake HANDSHAKE_TIMEOUT (10 s), armed once by the first touch Expired::Handshake; the core returns handshake_timed_out(), an io::ErrorKind::TimedOut error
Sniff SNIFF_TIMEOUT (300 ms), armed by enter(Phase::Sniff) Expired::Sniff; the core opens the flow with what it collected
Relay, Closing RELAY_IDLE_TIMEOUT (300 s), re-armed by touch on every byte event Expired::Idle; Timing has already pushed Effect::Finish and moved to Closing

Re-arming on every byte event is what makes the relay limit an idle limit rather than a lifetime. Timing::new arms nothing: a runtime delivers no event before the client speaks, so a silent client is the application’s watchdog to catch. For example, drive in app/src/serve.rs and serve_stream in protocols/src/tun/inbound.rs both run the runtime in its showing_progress stream mode and wrap every step taken before is_established() turns true in tokio::time::timeout(HANDSHAKE_TIMEOUT, runtime.next()).

SniffPrefix wraps the sniff Collector. push(plain) takes at most the remaining SNIFF_LIMIT (4 KiB) budget and returns (taken, Verdict), where Verdict is Found, More or Exhausted. The core consumes exactly taken. held() exposes the collected bytes for the ForwardHeld, and result() returns the SniffedBehavior as an Option. clear() drops the bytes; a core calls it at the top of the next byte event, which the held-buffer rule makes safe.

Passthrough<K> is the relay tail of a protocol with one stream outbound. It tracks the two end-of-stream flags and pushes the matching half-close:

Method Effects pushed
on_transport(data, fx) Forward { key, range: 0..data.len() }; returns data.len()
on_outbound(data, fx) none; stages data verbatim, or fails with staging_full()
on_outbound_eof(fx) ShutdownTransport, then Finish if the transport already ended
on_transport_eof(fx) Shutdown { key }, then Finish if the outbound already ended
on_outbound_gone(fx) ShutdownTransport, Finish (after ConnectFailed or OutboundError)

PassthroughCore<T> relays a flow whose destination is known before the first byte. The TUN inbound uses it, where the IP stack has already said where the client is going. It declares Key = Single, Target = Flow<T>, TransportAddr = (), STAGING_RESERVE = 0 and BUF_SIZE = 8 * 1024.

protocols/src/core/mod.rs
impl<T> PassthroughCore<T> {
pub const BUF_SIZE: usize = 8 * 1024;
pub fn new(flow: Flow<T>) -> Self;
pub fn sniffing(flow: Flow<T>) -> Self;
pub fn is_established(&self) -> bool;
}

Traced through passthrough_core_opens_on_the_first_bytes_and_relays_verbatim and a_sniffing_passthrough_core_holds_the_prefix_and_opens_with_the_host:

  1. First bytes, no sniffing. Transport(b"hello") opens the flow at once, because a runtime delivers no event before the client’s first bytes. The effects are Open { key: Single, target }, SetDeadline(Some(RELAY_IDLE_TIMEOUT)) from enter(Phase::Relay), a second SetDeadline from touch, and Forward { range: 0..5 }. The call returns 5.

  2. Downlink. Outbound { data: b"world" } pushes SetDeadline(Some(RELAY_IDLE_TIMEOUT)) and stages world verbatim.

  3. Both halves close. TransportEof pushes Shutdown { key: Single }; OutboundEof then pushes ShutdownTransport and Finish.

  4. Sniffing instead. Built with PassthroughCore::sniffing and an IP destination (worth_sniffing), the first Transport(b"GET / HTTP/1.1\r\nHo") only pushes SetDeadline(Some(SNIFF_TIMEOUT)) and consumes everything into SniffPrefix.

  5. Host found. The next Transport(b"st: example.com\r\n\r\nbody") completes the Host header. The core pushes Open with sniffed.domain == "example.com", ForwardHeld { range: 0..41 } (all collected bytes, body included) and SetDeadline(Some(RELAY_IDLE_TIMEOUT)).

  6. Held buffer released. Under a runtime, the following Transport(b"more") arrives only after the held forward has been applied. The core re-arms the idle deadline (SetDeadline(Some(RELAY_IDLE_TIMEOUT))), clears the prefix, and forwards 0..4 from the event slice. The harness applies nothing, so the test checks only that held() is empty afterwards.

If the sniff deadline fires first, Expired::Sniff opens the flow with sniffed: None and a ForwardHeld of what was collected (a_sniffing_passthrough_core_opens_on_the_sniff_deadline). A domain destination skips sniffing entirely (a_sniffing_passthrough_core_with_a_domain_target_opens_at_once).

CoreHarness in protocols/src/core/harness.rs is a hand-driven stand-in for the runtime. It owns the core, an EffectList, a WriteBuffer<HARNESS_STAGING> (HARNESS_STAGING = 64 * 1024) and, for datagram cores, a PacketList. It does no I/O.

protocols/src/core/harness.rs
pub const HARNESS_STAGING: usize = 64 * 1024;
pub struct CoreHarness<C: ProxyCoreDecode> {
pub core: C,
/* private fields */
}
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>;
}
Method What it does
new / over_datagrams A harness for a byte-stream core, or one whose sink is built with with_packets so stage_to works.
event Deliver any event once; returns the consumed count and the effects pushed.
transport event(Event::Transport(data)).
feed Deliver transport bytes the way the runtime does: call again on the unconsumed tail while the core makes progress. Forward and SendTo ranges in the result are rebased to be absolute in data; held ranges are left alone.
outbound event(Event::Outbound { key, data }).
staged Take everything staged toward the transport so far. Staging accumulates across calls until you take it.
staged_packets Take the staged packets with their peers, in order.
held The core’s held bytes under range, as a forward would resolve them.

A typical unit test, from protocols/tests/unit/core/mod.rs:

protocols/tests/unit/core/mod.rs
let mut h = CoreHarness::new(PassthroughCore::sniffing(flow()));
let mut opaque = [0x16u8, 0x03, 0x01];
h.transport(&mut opaque).unwrap();
let (_, fx) = h.event(Event::Deadline).unwrap();
assert!(matches!(&fx[0], Effect::Open { target, .. } if target.sniffed.is_none()));
assert!(matches!(&fx[1], Effect::ForwardHeld { range, .. } if *range == (0..3)));
assert!(h.core.is_established());

protocols/tests/unit/core/mod.rs also uses a smaller pattern: a with_sink function that builds Effects::new over a WriteBuffer::<256>, runs a closure, and returns the pushed effects and the staged bytes. It tests the helpers that are not whole cores (Timing, Passthrough), and the two non-sniffing PassthroughCore tests call handle through it directly.

A core fails the connection by returning Err from handle. Outbound failures are not fatal to the runtime: they reach the core as events and the core decides. Everything else that ends a connection early is a RuntimeError:

concepts/src/runtime.rs
pub enum RuntimeError<E> {
Transport(io::Error),
Core(E),
UnknownKey,
DuplicateKey,
RangeOutOfBounds,
BadConsume,
FrameTooLarge,
WrongLinkKind,
StagedWithoutPeer,
}
Variant Cause Display text
Transport The transport read failed (including a datagram link that reports a stream read), or a byte-stream transport’s write, flush or shutdown failed or wrote zero bytes. A packet a datagram transport refuses is Event::TransportSendFailed, not this error. transport: …
Core handle returned Err proxy core: …
UnknownKey A forward, send or shutdown for a key with no slot effect targets an unknown outbound key
DuplicateKey Open on a live key open reuses a live outbound key
RangeOutOfBounds A range outside the event slice or the held buffer forward range outside the event slice or held buffer
BadConsume Consumed more than given, or less than a whole outbound payload or transport datagram core consumed an impossible byte count
FrameTooLarge The unparsed region fills BUF_SIZE and the core wants more protocol frame exceeds the read buffer
WrongLinkKind A stream effect on a datagram outbound, or the reverse stream effect on a datagram outbound or vice versa
StagedWithoutPeer Bytes staged toward a datagram transport outside a packet bytes staged toward a datagram transport without a peer

UnknownKey, DuplicateKey, RangeOutOfBounds, BadConsume, WrongLinkKind and StagedWithoutPeer are core bugs, not peer behaviour: a core turns malformed wire input into its own Error, never into an out-of-range effect. FrameTooLarge means either a peer sent a frame longer than the buffer or BUF_SIZE is too small for the protocol.

Cancellation is the runtime’s: dropping the ProxyServerRuntime future or stream closes the transport and every outbound with it. A core has no Drop obligations towards the I/O.

Constant Value Where
ProxyCoreDecode::MAX_DATAGRAM default 4096 concepts/src/core.rs
INLINE_EFFECTS 4 concepts/src/core.rs
HARNESS_STAGING 64 * 1024 protocols/src/core/harness.rs
HANDSHAKE_TIMEOUT 10 s protocols/src/core/mod.rs
RELAY_IDLE_TIMEOUT 300 s protocols/src/core/mod.rs
SNIFF_TIMEOUT 300 ms protocols/src/sniff/mod.rs
SNIFF_LIMIT 4 * 1024 protocols/src/sniff/mod.rs

The in-tree cores and the constants they declare:

Core Key TransportAddr STAGING_RESERVE MAX_DATAGRAM BUF_SIZE
PassthroughCore (protocols/src/core/mod.rs) Single () 0 default 8 * 1024
HttpCore (protocols/src/http/core.rs) Single () 256 default MAX_HEAD (64 * 1024)
ShadowsocksCore (protocols/src/ss_legacy/core.rs) Single () 32 + 2 * CHUNK_OVERHEAD + 28 default 20 * 1024
Ss2022Core (protocols/src/ss_2022/core.rs) Single () 32 + 1 + 8 + 32 + 2 + 2 * TAG_SIZE + RECORD_OVERHEAD + 115 default 32 * 1024
TrojanCore (protocols/src/trojan/core.rs) FlowKey () PACKET_HEADER_MAX.next_multiple_of(16) + downlink_overhead(Self::BUF_SIZE) MAX_LENGTH (8192) 16 * 1024
VlessCore (protocols/src/vless/core.rs) FlowKey () 272 + downlink_overhead(Self::BUF_SIZE) 8192 16 * 1024
VMessCore (protocols/src/vmess/core.rs) FlowKey () 4096 8192 32 * 1024
Hy2StreamCore (protocols/src/hysteria/server/inbound.rs) Single () 2048 default 8 * 1024
Hy2UdpCore (protocols/src/hysteria/server/datagrams.rs) u32 () 4096 MAX_UDP_SIZE (4096) 16 * 1024
TunUdpCore (protocols/src/tun/udp.rs) Single SocketAddr 4096 4096 8 * 1024

The cores with FlowKey carry mux.cool sub-flows; downlink_overhead in protocols/src/mux/demux.rs accounts for the Keep frame headers one outbound read may be split into.

Test File Pins
core_is_driven_without_any_io concepts/tests/runtime.rs A core runs with hand-built events; split frames stay unconsumed; in-place decryption
relays_two_keys_and_completes_on_fin concepts/tests/runtime.rs Forwards queued behind Open wait for the connect; Finish completes with Traffic totals
stalled_outbound_holds_uplink_but_not_other_downlink concepts/tests/runtime.rs A blocked forward holds everything behind it; other downlink keeps flowing
connect_failure_reaches_the_core_as_an_event concepts/tests/runtime.rs ConnectFailed delivery
half_close_propagates_both_ways_and_finishes concepts/tests/runtime.rs Shutdown, OutboundEof, ShutdownTransport
deadline_event_lets_the_core_time_out, deadline_is_armed_against_the_tokio_clock concepts/tests/runtime.rs The single deadline and the tokio clock
frame_larger_than_the_buffer_is_an_error concepts/tests/runtime.rs FrameTooLarge
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites, a_held_range_past_the_buffer_is_rejected concepts/tests/runtime.rs The held-buffer rule, RangeOutOfBounds
datagram_outbound_is_never_truncated_by_staging_backpressure concepts/tests/runtime.rs Room for a whole packet before a datagram read
datagram_transport_demultiplexes_peers_and_frames_replies, staging_without_a_peer_under_a_datagram_transport_is_an_error, stage_to_is_refused_under_a_stream_transport concepts/tests/runtime.rs Datagram transports and StagedWithoutPeer
a_refused_datagram_send_keeps_the_key_alive, a_refused_transport_packet_is_reported_not_fatal, single_peer_datagram_transport_demuxes_sessions_and_reports_refusals concepts/tests/runtime.rs SendFailed and TransportSendFailed are not fatal
stream_mode_reports_each_unit_of_work concepts/tests/runtime.rs Per-step Traffic deltas; non-payload steps are zero deltas
timing_arms_handshake_once_then_idle_per_byte_event, timing_reports_handshake_and_sniff_expiry_to_the_core protocols/tests/unit/core/mod.rs Timing
sniff_prefix_finds_a_host_across_pushes_and_keeps_the_bytes, sniff_prefix_takes_no_more_than_its_budget protocols/tests/unit/core/mod.rs SniffPrefix and SNIFF_LIMIT
passthrough_half_closes_each_side_and_finishes_on_the_second protocols/tests/unit/core/mod.rs Passthrough
passthrough_core_* (2 tests, through with_sink), a_sniffing_passthrough_core_* (3 tests, through CoreHarness) protocols/tests/unit/core/mod.rs PassthroughCore: opening, relaying, half-close, failure, idle expiry and sniffing
stream_sessions_open_forward_and_end_with_fresh_generations protocols/tests/unit/mux/demux.rs Generation-tagged SubKeys
a_frame_straddling_chunks_is_held_and_forwarded_from_the_held_buffer, vmess_keeps_every_frame_one_read_completes protocols/tests/unit/mux/demux.rs Mux frames split across VMess chunks are forwarded from the held buffer
a_chunk_split_across_reads_decodes_its_header_once protocols/tests/unit/vmess/core.rs Never decrypting a header twice (VMess)
chunk_encoder_and_decoder_agree_and_never_decrypt_twice protocols/tests/unit/ss_legacy/aead.rs Never decrypting a header twice (Shadowsocks ChunkDecoder)

Run them with cargo test -p etemenanki-concepts --test runtime and cargo test -p etemenanki-protocols --lib core::tests (the filter also picks up each protocol core’s own core::tests module).