Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,17 @@ jobs:
- uses: actions/checkout@v4
- uses: dtolnay/rust-toolchain@stable
with:
components: clippy
components: clippy,rustfmt
- uses: Swatinem/rust-cache@v2
- name: Build
run: cargo build --lib --tests --examples
- name: Test (debug)
run: cargo test --lib
- name: Test (release)
run: cargo test --lib --release
- name: Stress (release)
run: cargo test --lib --release test_soak_randomized_stream_lifecycles -- --ignored
- name: Format
run: cargo fmt --all -- --check
- name: Clippy
run: cargo clippy --lib --tests
run: cargo clippy --lib --tests --examples --all-features -- -D warnings
5 changes: 3 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,8 +87,9 @@ let (connector, acceptor, worker) = MuxBuilder::server()
// Per-stream idle timeout (seconds): close streams with no
// recent traffic.
.with_idle_timeout(NonZeroU64::new(60).unwrap())
// Backpressure thresholds: cap how many frames may sit in the
// tx/rx queues before poll_write / poll_read park.
// Backpressure thresholds: cap queued tx frames and the combined
// inbound-frame/unaccepted-stream backlog. Keep-alive expiry pauses
// while the RX budget deliberately prevents carrier reads.
.with_max_tx_queue(NonZeroUsize::new(1024).unwrap())
.with_max_rx_queue(NonZeroUsize::new(1024).unwrap())
.with_connection(connection)
Expand Down
13 changes: 7 additions & 6 deletions src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,23 +83,24 @@ impl MuxBuilder<WithConfig> {
}

/// Per-stream idle timeout: if a stream sees no traffic for this
/// many seconds, it is closed and its handle is reaped.
/// many seconds, it is closed.
pub fn with_idle_timeout(&mut self, timeout_secs: NonZeroU64) -> &mut Self {
self.state.config.idle_timeout = Some(timeout_secs);
self
}

/// Backpressure threshold for outbound frames. `poll_write` parks
/// once a stream's pending tx queue exceeds this value. Defaults
/// Backpressure threshold for outbound frames per stream. `poll_write` parks
/// once a stream's pending tx queue reaches this value. Defaults
/// to 1024.
pub fn with_max_tx_queue(&mut self, size: NonZeroUsize) -> &mut Self {
self.state.config.max_tx_queue = size;
self
}

/// Backpressure threshold for inbound frames. The dispatcher parks
/// once total pending rx exceeds this value, propagating
/// backpressure to the peer's tx side. Defaults to 1024.
/// Backpressure threshold for inbound frames and streams waiting to be
/// accepted. The dispatcher parks once the total reaches this value,
/// propagating backpressure to the peer's tx side. Keep-alive expiry is
/// suspended while reads are deliberately parked. Defaults to 1024.
pub fn with_max_rx_queue(&mut self, size: NonZeroUsize) -> &mut Self {
self.state.config.max_rx_queue = size;
self
Expand Down
4 changes: 3 additions & 1 deletion src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ pub struct MuxConfig {
pub idle_timeout: Option<NonZeroU64>,
/// Backpressure threshold for outbound frames per stream.
pub max_tx_queue: NonZeroUsize,
/// Backpressure threshold for inbound frames across the session.
/// Backpressure threshold for inbound frames and unaccepted streams
/// across the session. Dead-peer detection is suspended while this
/// budget is exhausted because the carrier is intentionally not polled.
pub max_rx_queue: NonZeroUsize,
}
2 changes: 2 additions & 0 deletions src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ pub enum MuxError {
InvalidControlFramePayload(u16),
#[error("Reserved stream id {0}")]
ReservedStreamId(u32),
#[error("NOP frame must use stream id 0, got {0}")]
InvalidNopStreamId(u32),
#[error("Duplicated stream id {0}")]
DuplicatedStreamId(u32),

Expand Down
27 changes: 19 additions & 8 deletions src/frame.rs
Original file line number Diff line number Diff line change
Expand Up @@ -95,25 +95,36 @@ impl Decoder for MuxCodec {
type Error = MuxError;

fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
src.reserve(HEADER_SIZE + MAX_PAYLOAD_SIZE + HEADER_SIZE);

if src.len() < HEADER_SIZE {
src.reserve(HEADER_SIZE - src.len());
return Ok(None);
}
let header = MuxFrameHeader::decode(src)?;
let len = header.length as usize;
if src.len() < HEADER_SIZE + len {
return Ok(None);
}
// Per smux v1, only PSH carries payload; SYN/FIN/NOP must be empty.
// Reject non-empty control frames so a peer cannot use them for
// bandwidth amplification or covert framing.
// These checks depend only on the header, so reject immediately rather
// than waiting for an invalid peer to deliver its claimed payload.
match header.command {
MuxCommand::Sync | MuxCommand::Finish | MuxCommand::Nop if header.length != 0 => {
return Err(MuxError::InvalidControlFramePayload(header.length));
}
_ => {}
}
match header.command {
MuxCommand::Nop if header.stream_id != 0 => {
return Err(MuxError::InvalidNopStreamId(header.stream_id));
}
MuxCommand::Sync | MuxCommand::Finish | MuxCommand::Push if header.stream_id == 0 => {
return Err(MuxError::ReservedStreamId(0));
}
_ => {}
}

let len = header.length as usize;
let frame_len = HEADER_SIZE + len;
if src.len() < frame_len {
src.reserve(frame_len - src.len());
return Ok(None);
}
src.advance(HEADER_SIZE);
let payload = src.split_to(len).freeze();

Expand Down
Loading
Loading