The server runtime
Source files: 26 · checked against Etemenanki 596916d · katana v3.0.1
Etemenanki/concepts/src/runtime.rsEtemenanki/concepts/src/buffer.rsEtemenanki/concepts/src/wake.rsEtemenanki/concepts/src/core.rsEtemenanki/concepts/src/link.rsEtemenanki/concepts/tests/runtime.rsEtemenanki/concepts/tests/client.rsEtemenanki/app/src/serve.rsEtemenanki/protocols/src/core/mod.rsEtemenanki/protocols/src/http/core.rsEtemenanki/protocols/src/http/protocol.rsEtemenanki/protocols/src/trojan/core.rsEtemenanki/protocols/src/trojan/protocol.rsEtemenanki/protocols/src/vless/core.rsEtemenanki/protocols/src/vmess/core.rsEtemenanki/protocols/src/mux/demux.rsEtemenanki/protocols/src/ss_legacy/core.rsEtemenanki/protocols/src/ss_2022/core.rsEtemenanki/protocols/src/hysteria/protocol.rsEtemenanki/protocols/src/hysteria/server/inbound.rsEtemenanki/protocols/src/hysteria/server/datagrams.rsEtemenanki/protocols/src/tun/inbound.rsEtemenanki/protocols/src/tun/udp.rsEtemenanki/protocols/tests/support/pipeline.rsEtemenanki/environment/tests/integration/udp.rskatana/src/serve.rs
ProxyServerRuntime drives one proxied connection. It is a single hand-written Future (or, on request, a Stream) that owns three things: the client-facing transport, every outbound the connection opens, and the sans-I/O protocol core. It moves bytes between them through three fixed-size buffers, polls only the outbounds that actually woke, and never spawns a task. Every server protocol except SOCKS runs its traffic inside these runtimes, so the rules on this page decide how HTTP, Trojan, VLESS, VMess, Shadowsocks, Shadowsocks 2022, Hysteria 2 and TUN connections behave under load.
This page is for contributors who change concepts/src/runtime.rs, concepts/src/buffer.rs or concepts/src/wake.rs, and for protocol authors who need to know exactly when their core is called and with how much room. The core’s side of the contract (ProxyCoreDecode, Event, Effect, the Effects sink) is on The server core. This page covers the driver that calls it.
Responsibilities
Section titled “Responsibilities”The runtime:
- reads the transport into a fixed read buffer and hands the unparsed bytes to the core;
- dials an outbound when the core asks (
Effect::Open), through theConnectorit was built with; - applies the core’s effects strictly in the order they were pushed, writing forwarded ranges to outbounds without copying them;
- reads outbounds into a scratch buffer and hands each chunk or packet to the core;
- writes whatever the core staged back to the transport;
- runs one deadline timer on the core’s behalf;
- counts the bytes it moves (
Traffic); - decides when not to read. All of its backpressure comes from that decision.
The runtime does not:
- parse or produce protocol bytes. That is the core’s job.
- choose destinations, route or resolve names. That happens inside the
Connectorand the links it returns. - keep timers of its own. The only timer is the one the core arms with
Effect::SetDeadline. - log. Outbound failures go to the core as events, and fatal errors go to the caller as a
RuntimeError.
Where it runs
Section titled “Where it runs”| Caller | Core | Constructor | Mode |
|---|---|---|---|
app/src/serve.rs → drive |
HttpCore, TrojanCore, VlessCore, VMessCore, ShadowsocksCore, Ss2022Core |
new over the accepted stream |
showing_progress |
katana src/serve.rs → drive |
TrojanCore, VlessCore, VMessCore, ShadowsocksCore, Ss2022Core |
new over the accepted stream |
showing_progress |
protocols/src/hysteria/server/inbound.rs → classifier |
Hy2StreamCore, one runtime per proxy stream |
new over QuicIo |
quiet |
protocols/src/hysteria/server/inbound.rs, per QUIC connection once it has authenticated, when UDP is enabled |
Hy2UdpCore, one runtime over the connection’s datagrams |
over_datagrams over QuicDatagrams |
quiet |
protocols/src/tun/inbound.rs → serve_stream |
PassthroughCore, one runtime per TCP stream |
new |
showing_progress |
protocols/src/tun/inbound.rs, per client source address |
TunUdpCore, one runtime over that source’s UDP flows |
over_datagrams over TunUdpLink |
quiet |
The SOCKS inbound has its own driver (socks.serve) and does not use this runtime.
Every production caller instantiates BUF_SIZE with the core’s own associated constant, for example ProxyServerRuntime::<{ Hy2UdpCore::<()>::BUF_SIZE }, _, _, _>::over_datagrams(...), or ProxyServerRuntime::<BUF, Core, _, _>::new(stream, core, connector) inside a drive that is generic over const BUF: usize:
| Core | BUF_SIZE |
STAGING_RESERVE |
MAX_DATAGRAM |
|---|---|---|---|
PassthroughCore (protocols/src/core/mod.rs) |
8 * 1024 |
0 |
default 4096 |
HttpCore |
MAX_HEAD = 64 * 1024 |
256 |
default 4096 |
TrojanCore |
16 * 1024 |
computed from frame overheads | MAX_LENGTH = 8192 |
VlessCore |
16 * 1024 |
computed from frame overheads | 8192 |
VMessCore |
32 * 1024 |
4096 |
8192 |
ShadowsocksCore |
20 * 1024 |
computed from frame overheads | default 4096 |
Ss2022Core |
32 * 1024 |
computed from frame overheads | default 4096 |
Hy2StreamCore |
8 * 1024 |
2048 |
default 4096 |
Hy2UdpCore |
16 * 1024 |
4096 |
MAX_UDP_SIZE = 4096 |
TunUdpCore |
8 * 1024 |
4096 |
4096 |
Key types
Section titled “Key types”ProxyServerRuntime
Section titled “ProxyServerRuntime”pub struct ProxyServerRuntime<const BUF_SIZE: usize, Core, Trans, Conn, Mode = ProxyRunsQuiet>where Core: ProxyCoreDecode, Conn: Connector<Core::Target>,{ /* private fields */ }
pub struct ProxyRunsQuiet;pub struct ProxyShowsProgress;| Parameter | Bound | Meaning |
|---|---|---|
BUF_SIZE |
const usize, must exceed Core::STAGING_RESERVE |
Size of each of the three buffers (transport read, transport staging, outbound scratch). It also caps the largest protocol frame the connection accepts. |
Core |
ProxyCoreDecode |
The protocol state machine. Its associated types fix the outbound key (Key), what gets dialed (Target), its error (Error) and the transport peer address (TransportAddr). |
Trans |
Transport<Addr = Core::TransportAddr> on the impls that drive it |
StreamTransport<T> (from new) or DatagramTransport<D> (from over_datagrams). Callers never name it. |
Conn |
Connector<Core::Target>, plus Conn::Datagram: DatagramLink<Addr = Destination> on the impls that drive it |
Dials a target into an Outbound::Stream or an Outbound::Datagram. |
Mode |
ProxyRunsQuiet (default) or ProxyShowsProgress |
Selects the Future impl, which resolves to the total Traffic, or the Stream impl, which yields one Traffic delta per unit of work. |
Constructors
Section titled “Constructors”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;}
impl<const BUF_SIZE: usize, Core, D, Conn> ProxyServerRuntime<BUF_SIZE, Core, DatagramTransport<D>, Conn, ProxyRunsQuiet>where Core: ProxyCoreDecode, D: DatagramLink<Addr = Core::TransportAddr>, Conn: Connector<Core::Target>,{ pub fn over_datagrams(link: D, core: Core, connector: Conn) -> Self;}
impl<const BUF_SIZE: usize, Core, Trans, Conn> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>where Core: ProxyCoreDecode, Conn: Connector<Core::Target>,{ pub fn showing_progress( self, ) -> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>;}-
newwraps the stream inStreamTransport, andover_datagramswraps the link inDatagramTransport. Both then call a privatebuild, which allocates the three buffers and setsprefer_transport = true, so the very first poll tries the transport before the outbounds. -
buildstarts with:assert!(BUF_SIZE > Core::STAGING_RESERVE,"BUF_SIZE must exceed the core's STAGING_RESERVE or nothing is ever read");This is a run-time
assert!, so a bad instantiation panics when the runtime is constructed, not when it is compiled. The reason: a stream outbound is read only withSTAGING_RESERVE + 1bytes of staging room, which a buffer no larger than the reserve never has, anddatagram_limitcomputesBUF_SIZE - Core::STAGING_RESERVE. -
showing_progressconsumes the quiet runtime and moves every field into theProxyShowsProgressform. It is only available on the quiet type, so a runtime changes mode at most once. Every caller does it right after construction.
Accessors
Section titled “Accessors”impl<const BUF_SIZE: usize, Core, Trans, Conn, Mode> ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, Mode>where Core: ProxyCoreDecode, Conn: Connector<Core::Target>,{ pub fn core(&self) -> &Core; pub fn summary(&self) -> Traffic;}
// On the impl that drives the runtime// (Trans: Transport<Addr = Core::TransportAddr>, Conn::Datagram: DatagramLink<Addr = Destination>):pub const fn datagram_limit() -> usize;core()lets the caller inspect the core between steps. The serving layers use it to ask whether the handshake is over (is_established()).summary()returns the bytes moved so far.datagram_limit()is the largest packet delivered whole from a datagram outbound:Core::MAX_DATAGRAM, capped atBUF_SIZE - Core::STAGING_RESERVE. ForTunUdpCorethat ismin(4096, 8192 - 4096)= 4096 bytes.- The type implements
UnpinwheneverTrans: Unpin. The transport isUnpin, and the connect futures and the timer are boxed, so nothing inside relies on being pinned.
Transports
Section titled “Transports”The client side is seen through one trait, so the runtime is written once for both kinds of transport. The trait is public, but a caller never needs to name it.
pub enum Received<A> { Bytes, Datagram(A), Eof,}
pub trait Transport: Unpin { const DATAGRAM: bool; type Addr;
fn poll_recv( &mut self, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>, ) -> Poll<io::Result<Received<Self::Addr>>>;
fn poll_send( &mut self, cx: &mut Context<'_>, data: &[u8], to: Option<&Self::Addr>, ) -> Poll<io::Result<usize>>;
fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;
fn poll_shutdown(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>>;}
pub struct StreamTransport<T>(pub T);pub struct DatagramTransport<D>(pub D);StreamTransport<T> |
DatagramTransport<D> |
|
|---|---|---|
| Bound | T: AsyncRead + AsyncWrite + Unpin |
D: DatagramLink |
DATAGRAM |
false |
true |
Addr |
() |
D::Addr: a SocketAddr for a link whose packets name a peer (a UDP socket serving many peers, or TUN’s per-source TunUdpLink, addressed by the far end), () for a single-peer link such as a QUIC connection’s datagrams |
poll_recv |
A read that fills nothing is Received::Eof; otherwise Received::Bytes |
Always Received::Datagram(from). There is no end of stream. |
poll_send |
poll_write; to is ignored |
poll_send_to(to). to = None is an InvalidInput error, “a datagram transport needs a peer per packet”. |
poll_flush, poll_shutdown |
Forwarded to T |
No-ops that return Ready(Ok(())) |
The outbound side uses the traits in concepts/src/link.rs:
pub trait Connector<Target> { type Stream: AsyncRead + AsyncWrite + Unpin; type Datagram: DatagramLink; type Future: Future<Output = io::Result<Outbound<Self::Stream, Self::Datagram>>>;
fn connect(&mut self, target: Target) -> Self::Future;}
pub enum Outbound<S, D> { Stream(S), Datagram(D),}Any FnMut(Target) -> Fut whose future yields an Outbound is a Connector, which is how the tests pass closures. See Links and types for DatagramLink, UdpOutbound and Destination.
Buffers
Section titled “Buffers”The connection allocates three boxed arrays of BUF_SIZE bytes once, in build, and never grows them. buffer::boxed_array builds each one on the heap directly, so a large BUF_SIZE never passes through the stack. Every copy a proxied byte goes through lands in one of these buffers, and backpressure comes from them being full.
flowchart LR client["client transport"] up["up: ReadBuffer"] core["Core::handle"] queue["effect queue"] out["outbound links"] scratch["scratch"] staging["staging: WriteBuffer"] client -->|"poll_recv"| up up -->|"Event::Transport"| core core -->|"Forward ranges"| queue queue -->|"poll_write, poll_send_to"| out out -->|"poll_read, poll_recv_from"| scratch scratch -->|"Event::Outbound, Event::Datagram"| core core -->|"stage, reserve, stage_to"| staging staging -->|"poll_send"| client
Uplink bytes are not copied by the runtime at all. The core decrypts them in place in the read buffer and forwards ranges of it, which the runtime writes to the outbound straight from up. Downlink bytes are read into scratch and copied once more, when the core seals them into staging.
The read buffer: ReadBuffer<BUF_SIZE>
Section titled “The read buffer: ReadBuffer<BUF_SIZE>”The transport is read into up, which is split into three regions:
0 ............... start ............... end .............. BUF_SIZE| parsed | unparsed | free || (a queued | (the tail of a | (the next read || forward may | split frame) | lands here) || still read it) | | || Method | Signature | Used for |
|---|---|---|
start |
pub fn start(&self) -> usize |
The base that forward ranges from a transport event are made absolute against |
unparsed |
pub fn unparsed(&mut self) -> &mut [u8] |
The slice handed to the core as Event::Transport, mutable so the core can decrypt in place |
unparsed_len |
pub fn unparsed_len(&self) -> usize |
Whether there is anything left to feed |
slice |
pub fn slice(&self, range: Range<usize>) -> &[u8] |
Resolving a queued forward’s absolute range when it is applied |
free_len, free |
pub fn free_len(&self) -> usize, pub fn free(&mut self) -> &mut [u8] |
The tail the next read goes into |
advance_end |
pub fn advance_end(&mut self, n: usize) |
Accounting for a read of n bytes |
advance_start |
pub fn advance_start(&mut self, n: usize) |
Accounting for what the core consumed |
compact |
pub fn compact(&mut self) |
Moving the unparsed tail to offset 0. A no-op when start == 0, and a reset when nothing is unparsed. |
is_saturated |
pub fn is_saturated(&self) -> bool |
start == 0 && end == N: the unparsed region fills the whole array |
Queued forwards hold absolute offsets into this array, so the parsed region cannot move while one of them is queued. Compaction is therefore lazy and guarded:
- Stream transport.
poll_transport_readcompacts only when the free tail is empty (free_len() == 0). If the buffer is saturated at that point, the core still wants more of a frame that cannot fit, and the connection fails withRuntimeError::FrameTooLarge. The read path returns early while a transport-sourced forward is queued, so compaction never moves bytes that a forward still references (see Source pins). - Datagram transport. The buffer is always empty before a read, because a packet is consumed whole and a forward still reading it blocks the next read.
poll_transport_recv_packettherefore compacts (resets) before every read, with adebug_assert_eq!that nothing is unparsed.
The staging buffer: WriteBuffer<BUF_SIZE>
Section titled “The staging buffer: WriteBuffer<BUF_SIZE>”Bytes bound for the transport are staged in staging. The core appends at end through a Staging writer, and the runtime writes [start .. end) out and slides start forward.
0 ............. start ............. end ............. BUF_SIZE| written | pending | room |impl<const N: usize> WriteBuffer<N> { pub fn is_empty(&self) -> bool; pub fn pending(&self) -> &[u8]; pub fn advance_start(&mut self, n: usize); pub fn room(&mut self) -> usize; pub fn staging(&mut self) -> Staging<'_>;}
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<()>;}advance_startresets both offsets to 0 once everything pending has been written, so a drained buffer always offers its full size.room()slides the pending bytes to the front when that is what makes room, which is whenstart != 0and the tail is full (end == N). This is the only copy the buffer makes on its own.room()takes&mut selffor that reason, andstaging()calls it first, so the core always gets the largest tail available.Staging::reserveclaimslenbytes at the tail and returns them for the core to fill in place, or returnsNonewithout claiming anything.putisreserveplus a copy. A frame is sealed where it lies, with no intermediateVec.
The scratch buffer
Section titled “The scratch buffer”scratch: Box<[u8; BUF_SIZE]> receives outbound reads. service_key reads at most:
(staging.room() - Core::STAGING_RESERVE).min(BUF_SIZE)bytes from a stream outbound;datagram_limit()bytes from a datagram outbound.
The first cap is what makes the event contract hold: an outbound chunk of n bytes is only ever delivered with at least STAGING_RESERVE + n bytes of staging free, so the core can seal all of it.
The rest of the state
Section titled “The rest of the state”| Field | Type | Purpose |
|---|---|---|
outbounds |
BTreeMap<Core::Key, Slot<...>> |
One slot per live outbound key |
ready |
ReadyQueue<Core::Key> |
Keys whose outbound signalled readiness |
drained |
VecDeque<Core::Key> |
Keys popped from ready and not yet serviced, oldest first |
starved |
Vec<Core::Key> |
Keys that wanted servicing but were held back by staging room or a pinned buffer |
packets |
PacketList<Core::TransportAddr>, which is VecDeque<(usize, A)> |
Under a datagram transport, the packets in staging, front first |
effects |
EffectList<Core>, which is SmallVec<[Effect<C>; INLINE_EFFECTS]> |
The sink the core pushes into; empty between calls |
queue |
SmallVec<[Queued<Core>; INLINE_EFFECTS]> |
Effects not yet applied, front first |
deadline, deadline_armed |
Option<Pin<Box<Sleep>>>, bool |
The single timer |
delta, summary |
Traffic |
Bytes moved in the current unit of work, and in total |
The flags up_needs_feed, transport_read_closed, transport_write_closed, transport_needs_flush, shutdown_transport, finished, failed and prefer_transport complete the state. Beyond the three buffers, opening an outbound costs one Arc<KeyWaker> and one boxed dial future, and the timer is boxed once, on the first SetDeadline(Some(_)).
Outbound slots
Section titled “Outbound slots”Each live key owns a Slot in outbounds:
enum LinkState<F, S, D> { Connecting(Pin<Box<F>>), Stream(S), Datagram(D),}
struct Slot<K, F, S, D> { link: LinkState<F, S, D>, waker: Arc<KeyWaker<K>>, read_closed: bool, write_closed: bool, need_flush: bool,}need_flush is set when a forward was written but the link’s poll_flush returned Pending. service_key finishes that flush before it reads the key again.
stateDiagram-v2 [*] --> Connecting: Open effect reaches the queue front Connecting --> Stream: dial returns Outbound Stream Connecting --> Datagram: dial returns Outbound Datagram Connecting --> [*]: dial fails, ConnectFailed Connecting --> [*]: Close effect Stream --> [*]: Close effect Stream --> [*]: write, flush, shutdown or read error, OutboundError Stream --> [*]: shut down and read to EOF Datagram --> [*]: Close effect Datagram --> [*]: receive error, OutboundError
- Open.
Effect::Openfails the connection withDuplicateKeyif the key already has a slot, including one that is still connecting. Otherwise the runtime mints a waker withready.waker(key), callsconnector.connect(target), stores the future boxed asConnecting, and callsschedule()so the dial is polled once. The dial starts when theOpenreaches the front of the queue, not when the core pushes it: anOpenqueued behind a stalled forward waits. - Connect. When the dial resolves, the slot becomes
StreamorDatagram, the key is scheduled again, and the core receivesEvent::Connected { key }. A failed dial removes the key (seeforget_keyunder Applying effects) and deliversEvent::ConnectFailed { key, error }. - Half-close. A stream slot is removed when both halves are closed.
Effect::Shutdownsetswrite_closedoncepoll_shutdowncompletes; a zero-byte read setsread_closedand deliversEvent::OutboundEof. Whichever happens second removes the slot. - Datagram slots.
Effect::Shutdownon a datagram slot completes at once and only marks the write side. A datagram link has no end of stream, so the slot ends throughCloseor a receive error. - Close.
Effect::Closeremoves the slot at once, in any state, and drops the link or the in-flight dial future. It does not purge effects queued after it for the same key, so a later forward to that key fails withUnknownKeyunless anOpenfor the key comes first.
Per-key wake-ups
Section titled “Per-key wake-ups”One task polls many outbounds. Polling all of them on every wake-up would cost O(N) per wake-up, so each outbound is polled with its own waker that records which key became ready. It is the scheme FuturesUnordered uses, without its intrusive list.
struct Shared<K> { ready: Mutex<Vec<K>>, parent: Mutex<Option<Waker>>,}
pub struct ReadyQueue<K> { shared: Arc<Shared<K>>,}
impl<K> ReadyQueue<K> { pub fn new() -> Self; pub fn register(&self, waker: &Waker); pub fn waker(&self, key: K) -> Arc<KeyWaker<K>>; pub fn drain_into(&self, into: &mut VecDeque<K>); pub fn is_empty(&self) -> bool;}
pub struct KeyWaker<K> { key: K, queued: AtomicBool, shared: Arc<Shared<K>>,}
impl<K: Copy> KeyWaker<K> { pub fn key(&self) -> K; pub fn clear_queued(&self);}
impl<K: Copy + Send + Sync + 'static> KeyWaker<K> { pub fn waker(self: &Arc<Self>) -> Waker; pub fn schedule(self: &Arc<Self>);}
impl<K: Copy + Send + Sync + 'static> Wake for KeyWaker<K> { /* wake, wake_by_ref */ }Both mutexes are parking_lot::Mutex. Because KeyWaker implements std::task::Wake, turning it into a Waker is a reference-count increment, not an allocation; the one allocation is the Arc minted per open.
sequenceDiagram participant IO as outbound I/O driver participant KW as KeyWaker of key participant RQ as ReadyQueue participant T as runtime task T->>RQ: register task waker, every poll_once IO->>KW: wake_by_ref KW->>KW: queued.swap(true) KW->>RQ: push key, only if it was not queued KW->>RQ: take the parent waker KW->>T: wake the parent T->>RQ: drain_into(drained), oldest first T->>KW: clear_queued T->>IO: poll the outbound with the key's own waker
| Rule | Mechanism |
|---|---|
| A key is queued at most once, however often it wakes | wake_by_ref pushes the key only if queued.swap(true, Ordering::AcqRel) returned false |
| A burst of wake-ups wakes the task once | wake_by_ref take()s the parent waker. poll_once calls register on every poll to put it back, and register skips the clone when the stored waker will_wake the new one. |
| A wake-up that lands during a poll is not lost | service_key calls clear_queued() (a Release store) before polling the outbound, so a wake-up during the poll queues the key again instead of being swallowed by the dedupe |
| A new key is polled once | waker(key) returns an unqueued waker, so the runtime calls schedule() after Open |
| A key that needs a poll without an I/O event gets one | schedule() is wake_by_ref(). The runtime calls it after Open, after a dial completes, after every read that returned data or a packet (readiness wakers only fire after a Pending), and when a starved key is released. Because it also wakes the parent, a key scheduled during a poll is guaranteed another poll. |
poll_outbounds refills drained from ready only when drained is empty. It then services keys front to back and returns after the first one that made progress. A key whose read succeeded is rescheduled to the back of the ready queue, so busy outbounds take turns. A key that returns Pending is dropped from drained; its own waker queues it again when its I/O is ready. A key that no longer has a slot (a stale wake-up after Close) is skipped.
The effect queue
Section titled “The effect queue”The core pushes effects into the effects sink during handle. After every call the runtime validates them and moves them onto queue, tagging each with the buffer its range indexes:
enum Source { Transport, Scratch, Held,}
struct Queued<C: ProxyCoreDecode> { effect: Effect<C>, source: Source,}enqueue(source, base, limit) checks every range and makes it absolute:
ForwardandSendToranges must satisfystart <= end <= limit, wherelimitis the length of the event’s byte slice. They are then shifted bybase:up.start()for a transport event, 0 for an outbound event. Events that carry no bytes are enqueued withlimit = 0, so they can only forward empty ranges.ForwardHeldandSendToHeldranges must satisfystart <= end <= core.held().len(). They stay relative and are resolved againstheld()again when applied, so aVecthat reallocates between the two is fine.- Any violation fails the connection with
RangeOutOfBounds.
Source pins
Section titled “Source pins”A queued ranged effect pins the buffer it reads: nothing new is read into that buffer until the effect has been applied. Queued::pins(source) answers the question per effect, and pinned(source) asks it of the whole queue.
| Queued effect | Pins | While it is queued |
|---|---|---|
Forward or SendTo from Event::Transport or Event::TransportDatagram |
Source::Transport |
The transport is not read, and the core is not fed more transport bytes (up_needs_feed is set instead) |
Forward or SendTo from Event::Outbound or Event::Datagram |
Source::Scratch |
No connected outbound is serviced; keys that wake are parked in starved |
ForwardHeld or SendToHeld |
Source::Held |
The transport is not read or fed, and no connected outbound is serviced |
Open, Shutdown, Close, ShutdownTransport, SetDeadline, Finish |
nothing | No restriction |
The held pin is what lets a core rewrite its own held buffer safely. While any held effect is queued (held_free() is false), the core receives no byte event: no Transport, TransportDatagram, Outbound, Datagram, OutboundEof or TransportEof. A slot that is still Connecting is exempt from the scratch and held pins, because polling a dial touches neither buffer, so Connected and ConnectFailed still arrive; so do Deadline, TransportSendFailed and the failure events raised while effects are applied. The pin starts only when the core’s call returns and its effects are enqueued. Within one call, a held range pushed earlier still reads held() after the call, so the core must leave those bytes in place until its next byte event. protocols/src/mux/demux.rs → Demux::feed_chunks follows this rule: VMessCore passes it every chunk that one read opened in a single call, and the demux trims its held buffer only at the start of the next read it is fed. The rules for core authors are on The server core.
Applying effects: drive_effects
Section titled “Applying effects: drive_effects”drive_effects applies the queue front to back and stops at the first effect that cannot complete. It returns whether it applied anything. Outbound writes are polled with the key’s own waker, not the task’s, so a blocked forward wakes exactly its key.
| Effect | Applied as | Blocks on | Errors and events |
|---|---|---|---|
Forward, ForwardHeld |
An empty range is dropped. Otherwise poll_write of the range, then poll_flush. A partial write advances range.start and the effect stays at the front. |
A Connecting slot; poll_write returning Pending. A pending flush sets need_flush and does not block. |
No slot: UnknownKey. Datagram slot: WrongLinkKind. A write of 0 bytes (WriteZero), a write error or a flush error: the key is dropped and the core receives OutboundError. |
SendTo, SendToHeld |
One poll_send_to(data, &to) |
A Connecting slot; Pending |
No slot: UnknownKey. Stream slot: WrongLinkKind. A refused send: the packet is dropped, the key stays live, and the core receives SendFailed. |
Open |
See Outbound slots | Never | Live key: DuplicateKey |
Shutdown |
poll_shutdown on a stream; immediate on a datagram slot |
A Connecting slot; Pending |
No slot: UnknownKey. I/O error: the key is dropped and the core receives OutboundError. |
Close |
Removes the slot | Never | None, even for an unknown key |
ShutdownTransport |
Sets shutdown_transport |
Never | None |
SetDeadline(after) |
Some: resets the boxed Sleep to tokio::time::Instant::now() + after (creating it the first time) and arms it. None: disarms it. |
Never | None |
Finish |
Sets finished |
Never | None |
When a key is lost, forget_key removes the slot and every queued Forward, SendTo, ForwardHeld, SendToHeld, Shutdown and Close aimed at it. A forward waiting on a dial that just failed has nowhere to go, and the core learns of the loss through exactly one event (ConnectFailed or OutboundError). Whatever the core answers is enqueued at the back of the queue, behind anything still waiting.
At the end of a pass that applied something, if keys are starved and neither the scratch nor the held buffer is pinned any more, the starved keys are rescheduled. The bytes written by forwards and sends are added to delta.outbound_tx.
Scheduling: one unit of work
Section titled “Scheduling: one unit of work”poll_once registers the task’s waker with the ReadyQueue and calls step, which performs at most one unit of work. step returns Some(true) for progress, Some(false) for nothing to do until woken, and None once the connection is over.
flowchart TB
poll["poll_once: ready.register(task waker)"]
dl{"armed deadline fired?"}
fx{"drive_effects applied an effect?"}
wr{"poll_transport_write made progress?"}
done{"finished, staging empty, shutdown done?"}
fin{"finished?"}
io["transport read and outbound service, preferred side first"]
moved{"one side made progress?"}
prog["progress: move delta into summary"]
complete["complete"]
pend["Pending"]
poll --> dl
dl -->|yes| prog
dl -->|no| fx
fx -->|yes| prog
fx -->|no| wr
wr -->|yes| prog
wr -->|no| done
done -->|yes| complete
done -->|no| fin
fin -->|yes| pend
fin -->|no| io
io --> moved
moved -->|"yes, prefer the side that did not move"| prog
moved -->|no| pend
In order:
- Deadline.
poll_deadlinepolls the armedSleep. The timer goes first so that a connection whose I/O is always ready cannot starve its own handshake or idle deadline. When it fires, the timer is disarmed and the core receivesEvent::Deadline. - Effects.
drive_effectsretries whatever is still queued, typically a forward that was waiting on an outbound or a dial. - Transport writes.
poll_transport_writewrites staged bytes: onepoll_sendper step over a stream, every sendable packet over a datagram transport. Once staging is empty it flushes, and then performs a requestedpoll_shutdown. - Completion check. The runtime completes once
finishedis set, staging is empty, and the transport’s write side is closed ifShutdownTransportwas requested. - Reads, alternating. Unless the core has finished,
steptriespoll_transport_readandpoll_outboundsin the orderprefer_transportgives. A transport read that made progress setsprefer_transport = false; an outbound that made progress sets it back totrue. Neither side can starve the other. - Nothing to do. The step reports no progress and
poll_oncereturnsPending. Every source polled has registered a waker by then: the timer and the transport with the task’sContext, and each outbound with its key waker, which leads back to the task through theReadyQueue.
A unit that made progress moves delta into summary and hands delta to the mode-specific impl.
WORK_BUDGET. The quiet Future::poll loops over units of work up to WORK_BUDGET (64) times. If the budget runs out while every unit still makes progress, it calls cx.waker().wake_by_ref() and returns Pending. One busy connection whose transport or outbounds are always synchronously ready therefore cannot monopolize an executor worker. The stream mode performs one unit per poll_next, so there the caller sets the pace.
Reading the transport
Section titled “Reading the transport”poll_transport_read over a stream transport:
- Returns without progress if the read side is closed,
Source::TransportorSource::Heldis pinned, orstaging.room() < Core::STAGING_RESERVE. - If
up_needs_feedis set, feeds the leftover unparsed bytes to the core first, and reports progress if that consumed anything or pinned the buffer again. - If the free tail is empty, fails with
FrameTooLargewhen the buffer is saturated, and compacts otherwise. - Reads into the free tail. End of stream sets
transport_read_closedand deliversEvent::TransportEof. Bytes are counted intransport_rxand fed.
feed_transport delivers Event::Transport(up.unparsed()) in a loop. After each call it checks the consumed count, enqueues the effects with base = up.start(), advances start, and applies the effects. It stops when nothing is unparsed, when the core consumed nothing, or when a pin appears (setting up_needs_feed). So the core is called again as long as it consumes something, and the runtime reads more only once the core has consumed everything or has stopped consuming. Staging room is checked once, before the read, not before each call in this loop.
Servicing an outbound key
Section titled “Servicing an outbound key”service_key(key) polls one outbound once:
- Returns if the key has no slot. Otherwise calls
clear_queued()and builds aContextfrom the key’s waker. - If
need_flushis set, finishes the flush first.Pendingreturns without progress; an error fails the key. - Computes the staging room this key needs:
STAGING_RESERVE + datagram_limit()for a datagram slot,STAGING_RESERVE + 1otherwise (including a slot that is still connecting). If the room is short, or the scratch or held buffer is pinned and the slot is notConnecting, pushes the key ontostarvedand returns. Connecting: polls the dial (see Outbound slots).Stream: unlessread_closed, reads intoscratch[..max]. A zero-byte read is end of stream. Otherwise the key is rescheduled,outbound_rxis counted, and the core receivesEvent::Outbound { key, data }, which it must consume whole.Datagram: receives one packet intoscratch[..datagram_limit()]. The key is rescheduled, and the core receivesEvent::Datagram { key, from, data }, which it must consume whole. A receive error fails the key; it is the only way a datagram outbound is lost withoutClose.
Backpressure
Section titled “Backpressure”The buffers never grow, so every backpressure rule has the same shape: the runtime declines to read something until there is somewhere for the result to go. Bytes that are not read stay in the kernel socket or the peer’s send window, and the peer’s own flow control does the rest.
When the transport is not read
Section titled “When the transport is not read”| Condition | Why |
|---|---|
transport_read_closed |
The transport already returned end of stream |
A transport-sourced forward is queued (pinned(Source::Transport)) |
Its bytes are still in the read buffer at absolute offsets; reading and compacting could move them. A stalled outbound therefore stalls the uplink instead of buffering it. |
A held effect is queued (!held_free()) |
The core must not receive a byte event while its held buffer is pinned |
staging.room() < Core::STAGING_RESERVE |
The core must be able to stage its reply. A client that stops reading fills staging, and then its own uplink stops too. |
finished |
step no longer reads once the core has finished |
When the runtime stops feeding because of a pin, it sets up_needs_feed, and the next read attempt feeds the leftover unparsed bytes before reading anything new.
Reserve checks on the outbound side
Section titled “Reserve checks on the outbound side”| Outbound | Serviced only with this much staging room | Read size |
|---|---|---|
| Stream (or still connecting) | STAGING_RESERVE + 1 |
room - STAGING_RESERVE, capped at BUF_SIZE |
| Datagram | STAGING_RESERVE + datagram_limit() |
datagram_limit() |
A datagram outbound waits for room for a whole packet, never less. A client that stops reading therefore stalls the downlink at the outbound socket instead of receiving truncated packets.
Starved keys
Section titled “Starved keys”A key that wakes while the runtime cannot take its data (step 3 of Servicing an outbound key) is pushed onto starved and not polled. Its queued flag is already cleared, so its own waker can still queue it. reschedule_starved drains starved and calls schedule() on every key that still has a slot. It runs:
- after a stream transport write accepted at least one byte, which freed staging room;
- after a datagram transport send loop sent or dropped at least one packet;
- at the end of a
drive_effectspass that applied something, once neither the scratch nor the held buffer is pinned any more.
A rescheduled key that is still short of room is simply starved again.
One stalled write holds the whole queue
Section titled “One stalled write holds the whole queue”Effects apply strictly in order, and a forward that cannot complete holds every effect behind it. For a core with one outbound, that is exactly the backpressure wanted. For a multiplexing core it means one slow sub-flow holds the uplink of all its siblings until it drains, because its transport-sourced forward pins the read buffer and the transport is not read. Downlink from the other outbounds keeps flowing: their reads land in scratch and are sealed straight into staging without waiting on the queue.
sequenceDiagram participant C as client transport participant R as runtime participant A as outbound 1, slow participant B as outbound 2 C->>R: DATA for key 1, 100 bytes R->>A: poll_write accepts 8 bytes, then Pending Note over R: Forward stays at the queue front, read buffer pinned C--)R: FIN, not read while the forward is queued B->>R: downlink bytes R->>C: sealed DATA for key 2 via staging A-->>R: key waker fires, writable again R->>A: remaining 92 bytes R->>C: read resumes, FIN parsed
This is the scenario stalled_outbound_holds_uplink_but_not_other_downlink in concepts/tests/runtime.rs pins.
Event contract
Section titled “Event contract”The runtime checks staging room before it reads, so that the core can seal what it is handed:
| Poll that leads to the event | Room checked first |
|---|---|
Transport read: Transport, TransportDatagram, TransportEof |
STAGING_RESERVE |
Stream outbound service: Outbound, OutboundEof, and Connected or ConnectFailed for a dial |
STAGING_RESERVE + 1. A chunk of n bytes is read only into room - STAGING_RESERVE, so STAGING_RESERVE + n is free when it arrives. |
Datagram outbound service: Datagram |
STAGING_RESERVE + datagram_limit() |
Deadline |
none |
Failures raised while effects are applied or packets are sent: OutboundError from a write, flush or shutdown, SendFailed, TransportSendFailed |
none |
Deadline is delivered regardless of room, so that an idle timeout can still fire against a peer that has stopped reading. The failure events in the last row are raised by whatever pass hit the failure, without a room check of their own. A core must tolerate stage or reserve returning None in those handlers. More generally, a core should turn None from the staging writer into an error it returns (the test cores write .ok_or("staging full")?), never treat it as unreachable.
Transport bytes are expected to be forwarded, not staged. The reserve is checked once before the transport is read, and feed_transport may then call the core several times on the same bytes, so only the first call is sure to find the reserve free. A core that echoes transport bytes back reads Staging::room and consumes fewer bytes when its answer does not fit.
The runtime checks the consumed count after every call:
| Event | The core must consume |
|---|---|
Transport(data) |
At most data.len(). The rest stays in the read buffer and is offered again, followed by newly read bytes. |
TransportDatagram, Outbound, Datagram |
Exactly data.len(). There is nowhere to keep a remainder. |
| Events without bytes | Anything; the return value is ignored |
A violation fails the connection with BadConsume.
Datagram transports
Section titled “Datagram transports”over_datagrams builds a runtime over a DatagramLink: a link whose packets name a peer (TransportAddr = SocketAddr), such as one UDP socket serving many peers or TUN’s TunUdpLink, or the datagram side of a QUIC connection (TransportAddr = ()). The core usually acts as a demultiplexer. Hy2UdpCore, for example, keys its outbounds by session id, opens a session’s outbound the first time it sees the id, and closes idle sessions from its deadline. TunUdpCore uses a single outbound key for the whole source, and stages each reply to the far-end address it came from, which TunUdpLink maps back to the client’s flow.
-
Reads follow the kernel socket rule: one packet per read, into the empty read buffer, with room for at most
Core::MAX_DATAGRAM.min(BUF_SIZE)bytes. A longer packet is truncated by the link’s receive. The packet is delivered asEvent::TransportDatagram { from, data }and must be consumed whole. If the link ever returnsReceived::BytesorReceived::Eof, the connection fails withRuntimeError::Transport(InvalidData, “datagram transport delivered a stream read”). -
Staging needs a peer for every byte.
deliverbuilds the sink withEffects::with_packets, and the core stages with:concepts/src/core.rs 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 transport_is_datagram(&self) -> bool;Each
stage_toreserveslenbytes and pushes(len, to)ontopackets. After the core returns,delivercompares how much the staging buffer grew with the sum of the new packet lengths. A byte staged with plainstagebreaks that equality, and the connection fails withStagedWithoutPeer. Under a stream transport it is the other way round:stage_toreturnsNoneandtransport_is_datagram()isfalse. -
Sends go packet by packet in
poll_transport_send_packets:poll_send(&pending[..len], Some(&to))for the front packet, then the packet is popped and staging advances bylen.Pendingstops the loop. A packet the link refuses is dropped, is not counted intransport_tx, and is reported to the core asEvent::TransportSendFailed { to, error }. One peer’s unreachable address must not take the others down, so a refused send is never fatal; fatal link errors surface on the read side. -
Shutdown. There is no write half to close. Once the packets have drained, a pending
ShutdownTransportonly marks the write side closed. -
End. A datagram transport never delivers
TransportEof, so such a runtime ends only throughEffect::Finish, aRuntimeError(a fatal read error from the link, for example), or being dropped.
Finish and completion
Section titled “Finish and completion”Effect::Finish sets finished. Because effects apply in order, everything the core pushed before Finish (a Shutdown of an outbound, for example) has already completed, or been dropped with a lost key, when it takes effect.
From then on, step no longer reads the transport or any outbound. It keeps applying what is still queued and writing staging to the transport. The runtime completes when staging is empty and, if ShutdownTransport was requested, the transport’s write side has been shut down. The quiet form resolves to Ok(summary), and the stream form ends with None. The transport and any outbounds still open close when the caller drops the completed runtime.
| Core pushes | Result |
|---|---|
ShutdownTransport, then Finish |
Staging drains, the transport is flushed and half-closed, and the runtime completes. The client reads end of stream. |
Finish alone |
Staging drains and the runtime completes. The runtime polls the transport’s flush once but does not wait for it, and the transport closes when the runtime is dropped. Push ShutdownTransport too when the last bytes must be flushed. |
| Nothing, after every source has closed | The runtime stays pending. A core must push Finish to end the connection, typically once the transport has reached end of stream and every outbound is gone. |
Traffic and the two modes
Section titled “Traffic and the two modes”pub struct Traffic { pub transport_rx: u64, pub transport_tx: u64, pub outbound_rx: u64, pub outbound_tx: u64,}
impl std::ops::AddAssign for Traffic { /* field-wise += */ }| Counter | Counted when |
|---|---|
transport_rx |
A stream read lands n bytes in the read buffer, or a transport packet of n bytes (after truncation) is received |
transport_tx |
A stream write accepts n bytes, or a transport packet is sent. Refused packets are not counted. |
outbound_rx |
An outbound read or packet of n bytes lands in scratch |
outbound_tx |
A forward writes n bytes, or a datagram send reports n bytes |
impl<const BUF_SIZE: usize, Core, Trans, Conn> Future for ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyRunsQuiet>where Core: ProxyCoreDecode, Trans: Transport<Addr = Core::TransportAddr>, Conn: Connector<Core::Target>, Conn::Datagram: DatagramLink<Addr = Destination>,{ type Output = Result<Traffic, RuntimeError<Core::Error>>; fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;}The future resolves once, to the total traffic over the connection’s lifetime. Each poll runs up to WORK_BUDGET units of work. The Hysteria 2 runtimes and the TUN UDP runtime use this form, awaiting it inside a task they spawn. The tests spawn it directly, for example tokio::spawn(ProxyServerRuntime::<256, _, _, _>::new(transport, core, connector)).
impl<const BUF_SIZE: usize, Core, Trans, Conn> Stream for ProxyServerRuntime<BUF_SIZE, Core, Trans, Conn, ProxyShowsProgress>where Core: ProxyCoreDecode, Trans: Transport<Addr = Core::TransportAddr>, Conn: Connector<Core::Target>, Conn::Datagram: DatagramLink<Addr = Destination>,{ type Item = Result<Traffic, RuntimeError<Core::Error>>; fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>>;}Each item is the Traffic delta of one unit of work. A delta with every counter zero is a non-payload step: a connect, a close, a deadline. The deltas add up to summary(). The stream must be polled for the connection to move.
The stream form exists so that a caller can act between steps. The serving loops in app/src/serve.rs → drive, katana’s src/serve.rs → drive and protocols/src/tun/inbound.rs → serve_stream all use it the same way: they poll runtime.next() under HANDSHAKE_TIMEOUT until runtime.core().is_established(), then keep polling until the stream ends. A silent client produces no runtime event, which is why the handshake watchdog lives in the caller. katana additionally races next() against its user-retirement signal and resets its PROGRESS_WATCHDOG timer on every non-zero delta. See Serving connections and katana serving.
Failure paths and cancellation
Section titled “Failure paths and cancellation”RuntimeError
Section titled “RuntimeError”pub enum RuntimeError<E> { Transport(io::Error), Core(E), UnknownKey, DuplicateKey, RangeOutOfBounds, BadConsume, FrameTooLarge, WrongLinkKind, StagedWithoutPeer,}
impl<E: fmt::Display> fmt::Display for RuntimeError<E>;impl<E: fmt::Debug + fmt::Display> std::error::Error for RuntimeError<E>;impl<E> From<io::Error> for RuntimeError<E>;A RuntimeError ends the connection. Apart from Transport, Core and FrameTooLarge, every variant means the core broke the runtime’s contract, so it points to a bug in a protocol core rather than to a misbehaving peer. WrongLinkKind can also mean that the connector returned the other kind of link from the one the core expected. From<io::Error> maps to Transport.
| Variant | Display |
Raised when |
|---|---|---|
Transport(e) |
transport: {e} |
A transport read, write, flush or shutdown returns an error; a stream write accepts 0 bytes (WriteZero); a datagram transport returns a stream read. A refused datagram send is not an error. |
Core(e) |
proxy core: {e} |
Core::handle returns Err, for any event |
UnknownKey |
effect targets an unknown outbound key |
A non-empty Forward or ForwardHeld, a SendTo or SendToHeld, or a Shutdown reaches the queue front for a key with no slot |
DuplicateKey |
open reuses a live outbound key |
Open reaches the queue front for a key that has a slot |
RangeOutOfBounds |
forward range outside the event slice or held buffer |
At enqueue, a range is reversed or ends past the event slice or held(); at apply, a held range no longer fits held() |
BadConsume |
core consumed an impossible byte count |
A Transport call reports more than it was given, or a TransportDatagram, Outbound or Datagram payload is not consumed whole |
FrameTooLarge |
protocol frame exceeds the read buffer |
The unparsed transport data fills BUF_SIZE and the core still wants more |
WrongLinkKind |
stream effect on a datagram outbound or vice versa |
Forward/ForwardHeld on a datagram slot, or SendTo/SendToHeld on a stream slot, found once the dial has completed |
StagedWithoutPeer |
bytes staged toward a datagram transport without a peer |
Under a datagram transport, the core staged bytes outside stage_to/put_to |
Outbound failures are events
Section titled “Outbound failures are events”Outbound failures never end the connection by themselves. They reach the core, and the core decides:
| Event | Raised when | Key afterwards |
|---|---|---|
ConnectFailed { key, error } |
The dial future resolves to Err |
Gone; queued effects for it dropped |
OutboundError { key, error } |
A write returns an error or 0 bytes, a flush or shutdown fails, a stream read fails, or a datagram receive fails | Gone; queued effects for it dropped |
SendFailed { key, to, error } |
A datagram outbound refuses one packet (for example UdpOutbound refusing a domain with Unsupported) |
Live |
TransportSendFailed { to, error } |
A datagram transport refuses one packet | Not applicable; the runtime carries on |
fail_outbound and report_send_failure run from inside drive_effects, so they only enqueue what the core answers; the pass that raised them, or the next step, applies it.
Stream mode after an error
Section titled “Stream mode after an error”poll_next yields the error once as Some(Err(e)), sets failed, and returns None on every later poll.
Cancellation by drop
Section titled “Cancellation by drop”Dropping the future or the stream cancels the connection. The runtime owns the transport, the core, every in-flight dial future and every outbound link, so all of them are dropped with it and their sockets close. Dropping does not flush: bytes still in staging are discarded and queued effects are never applied. There is no separate cancel method, and the runtime spawns nothing that could outlive it. Callers rely on this: katana’s drive returns out of its select! when a user is retired, and the handshake watchdogs return on timeout, and in both cases dropping the runtime is the whole of the cleanup.
Invariants
Section titled “Invariants”| Invariant | Enforced by | Pinned by |
|---|---|---|
BUF_SIZE leaves room beyond the staging reserve |
assert! in build |
none |
| Bytes a queued ranged effect references are not moved or overwritten | Source pins: pinned(Source::Transport) blocks transport reads, compaction and feeding; scratch and held pins park outbound keys in starved; held_free() blocks byte events |
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites, stalled_outbound_holds_uplink_but_not_other_downlink |
| Effects apply in the order pushed | drive_effects works front to back and stops at the first blocked effect |
half_close_propagates_both_ways_and_finishes, stalled_outbound_holds_uplink_but_not_other_downlink |
| A forward to a key still dialing waits for the dial | LinkState::Connecting breaks out of drive_effects |
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites |
| A protocol frame fits the read buffer, or the connection fails | is_saturated() check raises FrameTooLarge |
frame_larger_than_the_buffer_is_an_error |
| An outbound chunk can always be sealed toward the transport | Room check in service_key; stream reads capped at room - STAGING_RESERVE |
exercised by every relay test, for example relays_two_keys_and_completes_on_fin |
| A datagram outbound is never truncated by backpressure | Room check of STAGING_RESERVE + datagram_limit() |
datagram_outbound_is_never_truncated_by_staging_backpressure |
| Every byte staged toward a datagram transport has a peer | Staged-bytes versus framed-bytes check in deliver |
staging_without_a_peer_under_a_datagram_transport_is_an_error, stage_to_is_refused_under_a_stream_transport |
| One refused packet does not end the connection | poll_transport_send_packets drops and reports; drive_effects reports SendFailed and keeps the key |
a_refused_transport_packet_is_reported_not_fatal, a_refused_datagram_send_keeps_the_key_alive, single_peer_datagram_transport_demuxes_sessions_and_reports_refusals |
| A lost key is reported once and nothing more is sent to it | forget_key before the single event |
connect_failure_reaches_the_core_as_an_event covers the event; the purge of queued effects has no dedicated test |
| A wake-up during a poll is not lost | clear_queued() before polling; schedule() after every successful read |
clearing_queued_lets_the_key_be_queued_again in concepts/src/wake.rs |
| A burst of outbound wake-ups wakes the task once | Dedupe on queued; parent waker taken on wake |
wake_queues_key_once_and_wakes_parent in concepts/src/wake.rs |
| The deadline cannot be starved by ready I/O | Polled first in step |
none directly; deadline_event_lets_the_core_time_out covers firing |
| The deadline follows tokio’s clock | tokio::time::Instant::now() and sleep_until in set_deadline |
deadline_is_armed_against_the_tokio_clock |
| One busy connection cannot monopolize a worker | WORK_BUDGET loop and self-wake in Future::poll |
none |
| Deltas add up to the total | poll_once moves delta into summary on every progress step |
stream_mode_reports_each_unit_of_work |
Limits
Section titled “Limits”| Constant | Value | Defined in | Meaning |
|---|---|---|---|
BUF_SIZE |
Per instantiation; see Where it runs | ProxyServerRuntime const generic |
Size of each of the three buffers, and the largest protocol frame |
ProxyCoreDecode::STAGING_RESERVE |
Declared per core | concepts/src/core.rs |
Staging room guaranteed before the transport is read |
ProxyCoreDecode::MAX_DATAGRAM |
4096 by default |
concepts/src/core.rs |
Largest datagram delivered whole, on either side |
datagram_limit() |
min(MAX_DATAGRAM, BUF_SIZE - STAGING_RESERVE) |
concepts/src/runtime.rs |
Read size for a datagram outbound |
| Transport packet read size | MAX_DATAGRAM.min(BUF_SIZE) |
poll_transport_recv_packet |
Read size for a datagram transport |
WORK_BUDGET |
64 |
concepts/src/runtime.rs |
Units of work per quiet Future::poll before a self-wake |
INLINE_EFFECTS |
4 |
concepts/src/core.rs |
Effects held inline in EffectList and in the queue before spilling to the heap (a handshake is Open + Forward + SetDeadline) |
The fixed part of a connection’s memory is the three BUF_SIZE arrays. Each open outbound adds a map entry, its link and one Arc<KeyWaker>, plus a boxed dial future while it is connecting. The timer adds one boxed Sleep once it is first armed. The effect queue spills to the heap beyond INLINE_EFFECTS entries, and the core bounds what it keeps in held() itself.
The integration tests live in concepts/tests/runtime.rs. They drive the runtime with toy sans-I/O cores over in-memory duplex streams, loopback UDP sockets and an mpsc-backed datagram link:
| Toy core | Transport | Key | STAGING_RESERVE |
Wire |
|---|---|---|---|---|
TinyMux |
Stream | u8 |
HDR + 64 = 68 |
[kind:u8][key:u8][len:u16][payload], payload XOR 0x55. Client kinds: 1 open, 2 data, 3 close, 4 fin. Replies: 2 data, 5 connected, 6 connect failed or outbound error, 7 EOF. |
TinyUdp |
Stream | Single |
4 | [len:u16][port:u16][payload]; len == 0 finishes |
TinyHub |
Datagram, SocketAddr peers |
SocketAddr |
2 | [port:u16][payload]; port 0 finishes |
Sniffing |
Stream | Single |
0 | Holds the first 8 bytes as the target name and forwards them from held() |
TinySessions over TinyQuic |
Datagram, single peer (()) |
u32 |
6 | [session:u32][port:u16][payload]; session 0 finishes |
| Test | Behaviour it pins |
|---|---|
core_is_driven_without_any_io |
A core runs on hand-built events with no runtime: a split frame stays unconsumed, and payloads are decrypted in place |
relays_two_keys_and_completes_on_fin |
Two keys relay uplink and one relays downlink; Connected reaches the client for both; ShutdownTransport plus Finish completes; exact Traffic totals |
stalled_outbound_holds_uplink_but_not_other_downlink |
A partial write holds the uplink behind it, while another key’s downlink still flows |
connect_failure_reaches_the_core_as_an_event |
A failed dial reaches the core as ConnectFailed, not as a runtime error |
half_close_propagates_both_ways_and_finishes |
Pending data reaches the outbound before its shutdown; outbound EOF reaches the core; the runtime finishes |
deadline_event_lets_the_core_time_out |
SetDeadline fires Deadline under a paused clock |
frame_larger_than_the_buffer_is_an_error |
FrameTooLarge with BUF_SIZE = 128 |
stream_mode_reports_each_unit_of_work |
showing_progress yields deltas that sum to the expected totals, one of them carrying the forward, and zero deltas for meta steps |
datagrams_round_trip_through_a_real_udp_socket |
SendTo and Event::Datagram through a real UDP outbound |
datagram_outbound_is_never_truncated_by_staging_backpressure |
With the client not reading, 1000-byte replies wait at the socket instead of being truncated |
deadline_is_armed_against_the_tokio_clock |
The deadline uses tokio’s clock, not std’s |
datagram_transport_demultiplexes_peers_and_frames_replies |
over_datagrams with many peers; each peer gets only its own replies; exact Traffic totals |
staging_without_a_peer_under_a_datagram_transport_is_an_error |
StagedWithoutPeer |
stage_to_is_refused_under_a_stream_transport |
stage_to returns None under a stream transport, while stage still works |
held_bytes_are_forwarded_after_the_dial_and_survive_later_rewrites |
A held forward waits for the dial, and later transport bytes wait for it |
a_held_range_past_the_buffer_is_rejected |
RangeOutOfBounds for a held range |
a_refused_datagram_send_keeps_the_key_alive |
SendFailed is reported and the same key keeps relaying |
a_refused_transport_packet_is_reported_not_fatal |
TransportSendFailed is reported, and the refused packet is not counted in transport_tx |
single_peer_datagram_transport_demuxes_sessions_and_reports_refusals |
A single-peer datagram link (TransportAddr = ()): an oversized reply is refused and other sessions carry on |
Unit tests next to the code:
| File | Test | Behaviour it pins |
|---|---|---|
concepts/src/buffer.rs |
read_buffer_compacts_unparsed_tail_to_front |
compact moves the unparsed tail to offset 0 |
concepts/src/buffer.rs |
read_buffer_saturation_means_frame_too_large |
is_saturated holds only when the unparsed region fills the array |
concepts/src/buffer.rs |
write_buffer_slides_pending_when_full |
room() slides pending bytes when the tail is full; a full drain resets |
concepts/src/buffer.rs |
staging_refuses_over_reservation_without_partial_commit |
A refused reserve claims nothing |
concepts/src/wake.rs |
wake_queues_key_once_and_wakes_parent |
Dedupe, and one parent wake per registration |
concepts/src/wake.rs |
clearing_queued_lets_the_key_be_queued_again |
A key is queued again only after clear_queued |
Other tests that run a server runtime:
concepts/tests/client.rs, moduleserver_side: a server runtime whose outbound is a client runtime in the same task, inserver_runtime_relays_through_a_client_runtime_in_one_task,upstream_dial_failure_is_connect_failed_not_connectedandrefused_upstream_handshake_is_connect_failed_not_connected.environment/tests/integration/udp.rs→dual_stack_link_serves_a_proxy_runtime: a runtime over the environment crate’s dual-stack UDP link.protocols/tests/support/pipeline.rs→serve_runtime: the harness the protocol tests use to run one runtime per accepted loopback connection.
The variants UnknownKey, DuplicateKey, WrongLinkKind and BadConsume, and the WORK_BUDGET self-wake, have no dedicated test. Add one when you touch those paths. To test a core without the runtime, see Testing.