From b796c95ea6a031c551e37f9d5f100b10bc508129 Mon Sep 17 00:00:00 2001 From: Martin Algesten Date: Fri, 2 Oct 2026 10:11:18 +0200 Subject: [PATCH 1/4] Validate DTLS 1.2 datagram order and back off flight resends --- src/dtls12/engine.rs | 232 +++++++++++++++++++++++++++++++++---- src/dtls12/incoming.rs | 165 ++++++++++++++++++++++++++ tests/dtls12/reorder.rs | 107 +++++++++++++++++ tests/dtls12/retransmit.rs | 16 +++ 4 files changed, 495 insertions(+), 25 deletions(-) diff --git a/src/dtls12/engine.rs b/src/dtls12/engine.rs index b60fd547..621d48f0 100644 --- a/src/dtls12/engine.rs +++ b/src/dtls12/engine.rs @@ -89,6 +89,11 @@ pub struct Engine { /// Timeout for the current flight flight_timeout: Timeout, + /// Cooldown for duplicate-triggered resends of the current flight. + /// Disabled allows a resend; Unarmed and Armed suppress it. Expiry only + /// allows another duplicate response, including after periodic retries stop. + flight_dupe_timeout: Timeout, + /// Global timeout for the entire connect operation. connect_timeout: Timeout, @@ -166,6 +171,7 @@ impl Engine { flight_saved_records: Vec::new(), flight_backoff, flight_timeout: Timeout::Unarmed, + flight_dupe_timeout: Timeout::Disabled, connect_timeout: Timeout::Unarmed, release_app_data: false, peer_handshake_confirmed: false, @@ -283,7 +289,7 @@ impl Engine { // drive a resend. if let Some(dupe_seq) = maybe_dupe_seq { if dupe_seq < self.peer_handshake_seq_no && !self.peer_handshake_confirmed { - if let Err(error) = self.flight_resend("dupe triggers resend") { + if let Err(error) = self.flight_resend_on_dupe() { self.recycle_incoming(incoming); return Err(error); } @@ -403,6 +409,15 @@ impl Engine { let timeout = now + self.flight_backoff.rto(); self.flight_timeout = Timeout::Armed(timeout); } + match self.flight_dupe_timeout { + Timeout::Unarmed => { + self.flight_dupe_timeout = Timeout::Armed(now + self.flight_backoff.rto()); + } + Timeout::Armed(deadline) if now >= deadline => { + self.flight_dupe_timeout = Timeout::Disabled; + } + _ => {} + } // The connect timeout is the overall timeout for establishing the connection if let Timeout::Armed(connect_timeout) = self.connect_timeout { @@ -426,6 +441,7 @@ impl Engine { let timeout = now + self.flight_backoff.rto(); self.flight_timeout = Timeout::Armed(timeout); self.flight_resend("flight timeout")?; + self.flight_dupe_timeout = Timeout::Armed(timeout); } else { return Err(Error::Timeout(crate::TimeoutError::Handshake)); } @@ -525,31 +541,25 @@ impl Engine { } fn poll_timeout(&self, now: Instant) -> Instant { - // No timeouts, return a distant future - if self.connect_timeout == Timeout::Disabled && self.flight_timeout == Timeout::Disabled { - const DISTANT_FUTURE: Duration = Duration::from_secs(10 * 365 * 24 * 60 * 60); - return now + DISTANT_FUTURE; - } - - match (self.connect_timeout, self.flight_timeout) { - // Keep this before the `(Armed, _)` arms. Starting a new flight resets its timer to - // `Unarmed`, but leaves the overall connection timer armed. If that mixed state - // returned the connection deadline, the caller would not drive `handle_timeout` to - // arm the flight timer until the whole handshake expired, so the flight would never - // be retransmitted. Returning `now` requests that immediate drive; the next poll sees - // both concrete deadlines and can return the earlier one. - (Timeout::Unarmed, _) | (_, Timeout::Unarmed) => now, - (Timeout::Armed(c), Timeout::Armed(f)) => { - if c < f { - c - } else { - f - } - } - (Timeout::Armed(c), _) => c, - (_, Timeout::Armed(f)) => f, - _ => now, + let timeouts = [ + self.connect_timeout, + self.flight_timeout, + self.flight_dupe_timeout, + ]; + // Request an immediate handle_timeout(now) to arm pending timers with + // fresh caller time. An armed connection deadline must not hide them. + if timeouts.contains(&Timeout::Unarmed) { + return now; } + const DISTANT_FUTURE: Duration = Duration::from_secs(10 * 365 * 24 * 60 * 60); + timeouts + .into_iter() + .filter_map(|timeout| match timeout { + Timeout::Armed(deadline) => Some(deadline), + _ => None, + }) + .min() + .unwrap_or(now + DISTANT_FUTURE) } pub fn flight_begin(&mut self, flight_no: u8) { @@ -557,6 +567,7 @@ impl Engine { self.flight_backoff.reset(&mut self.rng); self.flight_clear_resends(); self.flight_timeout = Timeout::Unarmed; + self.flight_dupe_timeout = Timeout::Disabled; } pub fn flight_stop_resend_timers(&mut self) { @@ -572,6 +583,7 @@ impl Engine { // authenticated application data instead). if self.is_client { self.peer_handshake_confirmed = true; + self.flight_dupe_timeout = Timeout::Disabled; } } @@ -581,6 +593,22 @@ impl Engine { } } + fn flight_resend_on_dupe(&mut self) -> Result<(), Error> { + if self.flight_dupe_timeout != Timeout::Disabled { + return Ok(()); + } + self.flight_resend("dupe triggers resend")?; + self.flight_backoff.attempt(&mut self.rng); + self.flight_dupe_timeout = Timeout::Unarmed; + // Restart the regular flight timer too, so its old deadline cannot + // produce another resend immediately after this one. Keep final-flight + // periodic retries disabled; only the duplicate cooldown remains active. + if self.flight_timeout != Timeout::Disabled { + self.flight_timeout = Timeout::Unarmed; + } + Ok(()) + } + fn flight_resend(&mut self, reason: &str) -> Result<(), Error> { debug!("Resending flight due to {}", reason); @@ -1064,6 +1092,7 @@ impl Engine { self.queue_tx.clear(); self.flight_saved_records.clear(); self.flight_timeout = Timeout::Disabled; + self.flight_dupe_timeout = Timeout::Disabled; self.connect_timeout = Timeout::Disabled; } @@ -1396,6 +1425,7 @@ impl RecordHandler for Engine { // confirms separately at its own completion (flight_stop_resend_timers). if content_type == ContentType::ApplicationData { self.peer_handshake_confirmed = true; + self.flight_dupe_timeout = Timeout::Disabled; } } @@ -1429,6 +1459,158 @@ impl RecordHandler for Engine { mod tests { use super::*; + fn poll_packets(engine: &mut Engine, now: Instant) -> (usize, Instant) { + let mut packets = 0; + let mut output = [0; 2048]; + loop { + match engine.poll_output(&mut output, now) { + Output::Packet(_) => packets += 1, + Output::Timeout(deadline) => return (packets, deadline), + _ => panic!("unexpected engine output"), + } + } + } + + fn duplicate_handshake(seq: u16) -> Vec { + let mut packet = vec![22, 0xFE, 0xFD, 0, 0, 0, 0, 0, 0, 0, 1, 0, 12, 14, 0, 0, 0]; + packet.extend_from_slice(&seq.to_be_bytes()); + packet.extend_from_slice(&[0; 6]); + packet + } + + fn resend_engine(now: Instant) -> Engine { + let config = Config::builder() + .dangerously_set_rng_seed(42) + .build() + .expect("valid config"); + let mut engine = Engine::new(Arc::new(config), AuthMode::Psk); + engine.peer_handshake_seq_no = 3; + engine.flight_begin(4); + engine + .create_record(ContentType::ChangeCipherSpec, 0, true, |buf| buf.push(1)) + .expect("save a flight"); + assert_eq!(poll_packets(&mut engine, now).0, 1); + engine.handle_timeout(now).expect("arm flight timer"); + assert_eq!(poll_packets(&mut engine, now).0, 0); + engine + } + + #[test] + fn ordered_datagrams_can_arrive_out_of_order() { + let now = Instant::now(); + let mut engine = Engine::new(Arc::new(Config::default()), AuthMode::Psk); + engine.peer_handshake_seq_no = 2; + let later = [duplicate_handshake(3), duplicate_handshake(4)].concat(); + engine + .parse_packet(&later) + .expect("queue sequences 3 and 4 first"); + assert_eq!(poll_packets(&mut engine, now).0, 0); + assert!(!engine.has_complete_handshake(MessageType::ServerHelloDone)); + engine + .parse_packet(&duplicate_handshake(2)) + .expect("queue sequence 2 later"); + assert_eq!(poll_packets(&mut engine, now).0, 0); + + let mut buffer = engine.pop_buffer(); + for seq in 2..=4 { + let handshake = engine + .next_handshake(MessageType::ServerHelloDone, &mut buffer) + .expect("reassemble in message order") + .expect("expected sequence is now available"); + assert_eq!(handshake.header.message_seq, seq); + assert_eq!(poll_packets(&mut engine, now).0, 0); + } + engine.push_buffer(buffer); + } + + #[test] + fn duplicate_resend_restarts_backoff_and_suppresses_other_old_messages() { + let now = Instant::now(); + let mut engine = resend_engine(now); + let initial_deadline = poll_packets(&mut engine, now).1; + let resend_at = initial_deadline - Duration::from_millis(1); + engine.handle_timeout(resend_at).expect("advance time"); + assert_eq!(poll_packets(&mut engine, resend_at).0, 0); + + engine + .parse_packet(&duplicate_handshake(1)) + .expect("first duplicate"); + assert_eq!(poll_packets(&mut engine, resend_at), (1, resend_at)); + engine + .handle_timeout(resend_at) + .expect("arm resend backoff"); + let (packets, deadline) = poll_packets(&mut engine, resend_at); + assert_eq!(packets, 0); + assert!(deadline > initial_deadline + Duration::from_secs(1)); + + for seq in [1, 2, 1, 2] { + engine + .parse_packet(&duplicate_handshake(seq)) + .expect("further duplicate"); + assert_eq!(poll_packets(&mut engine, resend_at), (0, deadline)); + } + engine + .handle_timeout(initial_deadline) + .expect("old deadline is cancelled"); + assert_eq!(poll_packets(&mut engine, initial_deadline), (0, deadline)); + engine + .handle_timeout(deadline) + .expect("retry at new deadline"); + let (packets, next_deadline) = poll_packets(&mut engine, deadline); + assert_eq!(packets, 1); + assert!(next_deadline - deadline > deadline - resend_at); + engine + .parse_packet(&duplicate_handshake(2)) + .expect("duplicate after timer resend"); + assert_eq!(poll_packets(&mut engine, deadline), (0, next_deadline)); + + engine.flight_begin(6); + engine + .create_record(ContentType::ChangeCipherSpec, 0, true, |buf| buf.push(1)) + .expect("save new flight"); + assert_eq!(poll_packets(&mut engine, deadline).0, 1); + engine + .parse_packet(&duplicate_handshake(2)) + .expect("new flight duplicate"); + assert_eq!(poll_packets(&mut engine, deadline).0, 1); + } + + #[test] + fn final_flight_duplicates_back_off_without_periodic_retransmissions() { + let now = Instant::now(); + let mut engine = resend_engine(now); + engine.flight_stop_resend_timers(); + poll_packets(&mut engine, now); + + engine + .parse_packet(&duplicate_handshake(1)) + .expect("first final-flight duplicate"); + assert_eq!(poll_packets(&mut engine, now), (1, now)); + engine + .handle_timeout(now) + .expect("arm final-flight cooldown"); + let (packets, deadline) = poll_packets(&mut engine, now); + assert_eq!(packets, 0); + assert!(deadline < now + Duration::from_secs(10)); + engine + .parse_packet(&duplicate_handshake(2)) + .expect("further final-flight duplicate"); + assert_eq!(poll_packets(&mut engine, now), (0, deadline)); + + engine.handle_timeout(deadline).expect("expire cooldown"); + assert_eq!(poll_packets(&mut engine, deadline).0, 0); + engine + .parse_packet(&duplicate_handshake(2)) + .expect("retry after cooldown"); + assert_eq!(poll_packets(&mut engine, deadline), (1, deadline)); + engine + .handle_timeout(deadline) + .expect("arm longer cooldown"); + let (packets, next_deadline) = poll_packets(&mut engine, deadline); + assert_eq!(packets, 0); + assert!(next_deadline - deadline > deadline - now); + } + #[test] fn duplicate_datagram_recycles_its_pooled_buffer() { let mut engine = Engine::new(Arc::new(Config::default()), AuthMode::Psk); diff --git a/src/dtls12/incoming.rs b/src/dtls12/incoming.rs index 8d7b45ea..e2a7e81d 100644 --- a/src/dtls12/incoming.rs +++ b/src/dtls12/incoming.rs @@ -80,6 +80,18 @@ impl Records { return Err(error); } + // The receive queue preserves each datagram's internal order. Discard an + // unordered datagram before classifying controls or retaining any of it, + // so a subsequent clean retransmission can recover. + if !Self::is_ordered(&parsed_records, decrypt.is_peer_encryption_enabled()) { + for record in parsed_records { + decrypt.push_buffer(record.into_buffer()); + } + return Ok(Records { + records: ArrayVec::new(), + }); + } + let mut records = ArrayVec::new(); for record in parsed_records { if let Some(record) = decrypt.classify_record(record)? { @@ -129,6 +141,58 @@ impl Records { Ok(()) } + + fn is_ordered(records: &[Record], peer_encryption_enabled: bool) -> bool { + let mut previous_epoch = 0; + let mut previous_seq = None; + let mut seen_ccs = false; + + for record in records { + let header = record.record(); + let epoch = header.sequence.epoch; + // After peer encryption is enabled, late plaintext records are + // discarded separately and may trail authenticated application data. + if !peer_encryption_enabled && epoch < previous_epoch { + return false; + } + previous_epoch = epoch; + + if header.content_type == ContentType::ChangeCipherSpec { + seen_ccs = true; + } + if header.content_type != ContentType::Handshake { + continue; + } + if !peer_encryption_enabled && seen_ccs && epoch == 0 { + return false; + } + // A queued future-epoch Finished is still ciphertext until CCS has + // enabled decryption; its bytes cannot be treated as message_seq. + if epoch > 0 && !peer_encryption_enabled { + continue; + } + + // Check the full payload, including messages beyond the bounded + // parsed handshake list. Match its handling of malformed tails. + let mut fragment = header.fragment(record.buffer()); + while !fragment.is_empty() { + let Ok((body, handshake)) = Handshake::parse_header(fragment) else { + break; + }; + let Some(remaining) = body.get(handshake.fragment_length as usize..) else { + break; + }; + let seq = handshake.message_seq; + if previous_seq.is_some_and(|previous| seq < previous) { + return false; + } + // Equal sequences are normal for fragments of one message. + previous_seq = Some(seq); + fragment = remaining; + } + } + true + } } impl Deref for Records { @@ -504,6 +568,107 @@ mod tests { out } + fn build_handshake(seq: u16, offset: u32, body: &[u8]) -> Vec { + let mut out = vec![14]; // ServerHelloDone; only the fragment header is parsed here. + out.extend_from_slice(&4u32.to_be_bytes()[1..]); + out.extend_from_slice(&seq.to_be_bytes()); + out.extend_from_slice(&offset.to_be_bytes()[1..]); + out.extend_from_slice(&(body.len() as u32).to_be_bytes()[1..]); + out.extend_from_slice(body); + out + } + + #[test] + fn unordered_handshakes_discard_the_whole_datagram_before_classification() { + for separate_records in [false, true] { + let mut packet = build_record(ContentType::Alert, 0, 0, &[2, 40]); + let mut handshakes = Vec::new(); + for seq in [3, 1, 2] { + let handshake = build_handshake(seq, 0, &[0; 4]); + if separate_records { + packet.extend_from_slice(&build_record( + ContentType::Handshake, + 0, + u64::from(seq), + &handshake, + )); + } else { + handshakes.extend_from_slice(&handshake); + } + } + if !separate_records { + packet.extend_from_slice(&build_record(ContentType::Handshake, 0, 1, &handshakes)); + } + + let mut handler = TestHandler::default(); + let incoming = Incoming::parse_packet(&packet, &mut handler, None) + .expect("unordered datagrams are silently discarded"); + assert!(incoming.is_none()); + assert_eq!(handler.classify_calls, 0); + assert_eq!(handler.buffers_returned, handler.buffers_acquired); + } + } + + #[test] + fn unordered_handshakes_beyond_the_parsed_capacity_are_discarded() { + let mut handshakes = Vec::new(); + for seq in 0..8 { + handshakes.extend_from_slice(&build_handshake(seq, 0, &[0; 4])); + } + handshakes.extend_from_slice(&build_handshake(0, 0, &[0; 4])); + let packet = build_record(ContentType::Handshake, 0, 0, &handshakes); + let mut handler = TestHandler::default(); + assert!( + Incoming::parse_packet(&packet, &mut handler, None) + .expect("unordered datagrams are silently discarded") + .is_none() + ); + assert_eq!(handler.classify_calls, 0); + assert_eq!(handler.buffers_returned, handler.buffers_acquired); + } + + #[test] + fn ordered_datagrams_allow_fragments_and_arbitrary_starting_sequences() { + let mut packet = Vec::new(); + for (seq, offset, body) in [ + (4, 0, &[0; 2][..]), + (4, 2, &[0; 2][..]), + (5, 0, &[0; 4][..]), + ] { + packet.extend_from_slice(&build_record( + ContentType::Handshake, + 0, + u64::from(seq), + &build_handshake(seq, offset, body), + )); + } + let mut handler = TestHandler::default(); + let incoming = Incoming::parse_packet(&packet, &mut handler, None) + .expect("parse ordered fragments") + .expect("ordered fragments must remain available"); + assert_eq!(incoming.records().len(), 3); + } + + #[test] + fn unordered_ccs_and_epochs_discard_the_whole_datagram() { + let ccs = build_record(ContentType::ChangeCipherSpec, 0, 1, &[1]); + let finished = build_record(ContentType::Handshake, 1, 0, &[0; 16]); + let handshake = build_record( + ContentType::Handshake, + 0, + 0, + &build_handshake(4, 0, &[0; 4]), + ); + for packet in [[finished, ccs.clone()].concat(), [ccs, handshake].concat()] { + let mut handler = TestHandler::default(); + let incoming = Incoming::parse_packet(&packet, &mut handler, None) + .expect("unordered datagrams are silently discarded"); + assert!(incoming.is_none()); + assert_eq!(handler.classify_calls, 0); + assert_eq!(handler.buffers_returned, handler.buffers_acquired); + } + } + #[test] fn parse_packet_filters_control_records_after_packet_validation() { let mut packet = Vec::new(); diff --git a/tests/dtls12/reorder.rs b/tests/dtls12/reorder.rs index 5b75a37f..95923a50 100644 --- a/tests/dtls12/reorder.rs +++ b/tests/dtls12/reorder.rs @@ -7,6 +7,113 @@ use dimpl::Dtls; use crate::common::*; +#[test] +fn dtls12_discards_scrambled_datagram_and_accepts_clean_retransmission() { + use dimpl::{Config, DtlsCertificate}; + + use crate::ossl_helper::{DtlsCertOptions, OsslDtlsCert}; + + let certificate = || { + let cert = OsslDtlsCert::new(DtlsCertOptions::default()); + DtlsCertificate { + certificate: cert.x509.to_der().expect("certificate DER"), + private_key: cert.pkey.private_key_to_der().expect("private key DER"), + } + }; + let config = Arc::new( + Config::builder() + .use_server_cookie(false) + .build() + .expect("valid config"), + ); + let now = Instant::now(); + let mut client = Dtls::new_12(Arc::clone(&config), certificate(), now); + client.set_active(true); + let mut server = Dtls::new_12(config, certificate(), now); + server.set_active(false); + server.handle_timeout(now).expect("initialize server"); + assert!(drain_outputs(&mut server).packets.is_empty()); + + client.handle_timeout(now).expect("start ClientHello"); + let client_hello = drain_outputs(&mut client); + client.handle_timeout(now).expect("arm client flight timer"); + assert!(drain_outputs(&mut client).packets.is_empty()); + let mut server_flight = Vec::new(); + for packet in &client_hello.packets { + server.handle_packet(packet).expect("receive ClientHello"); + server_flight.extend(drain_outputs(&mut server).packets); + } + server.handle_timeout(now).expect("arm server flight timer"); + assert!(drain_outputs(&mut server).packets.is_empty()); + let mut client_flight = Vec::new(); + for packet in &server_flight { + client.handle_packet(packet).expect("receive server flight"); + client_flight.extend(drain_outputs(&mut client).packets); + } + assert_eq!( + client_flight.len(), + 1, + "reproduce reordering within one datagram" + ); + client + .handle_timeout(now) + .expect("arm client final-flight timer"); + let deadline = drain_outputs(&mut client) + .timeout + .expect("client retry deadline"); + + let mut remaining = &client_flight[0][..]; + let mut records = Vec::new(); + while !remaining.is_empty() { + let length = 13 + usize::from(u16::from_be_bytes([remaining[11], remaining[12]])); + records.push(&remaining[..length]); + remaining = &remaining[length..]; + } + assert!( + records.len() >= 5, + "client flight includes certificate, key exchange, verify, CCS and Finished" + ); + let scrambled: Vec = records + .iter() + .rev() + .flat_map(|record| record.iter().copied()) + .collect(); + server + .handle_packet(&scrambled) + .expect("silently discard scrambled datagram"); + let discarded = drain_outputs(&mut server); + assert!(discarded.packets.is_empty()); + assert!(!discarded.connected); + + client + .handle_timeout(deadline) + .expect("retransmit clean final flight"); + let retry = drain_outputs(&mut client); + assert!(!retry.packets.is_empty()); + let mut final_flight = Vec::new(); + let mut server_connected = false; + for packet in &retry.packets { + server + .handle_packet(packet) + .expect("accept clean retransmission"); + let output = drain_outputs(&mut server); + server_connected |= output.connected; + final_flight.extend(output.packets); + } + assert!( + server_connected, + "scrambled copy must not poison the receive queue" + ); + let mut client_connected = false; + for packet in &final_flight { + client + .handle_packet(packet) + .expect("receive server final flight"); + client_connected |= drain_outputs(&mut client).connected; + } + assert!(client_connected); +} + #[test] #[cfg(feature = "rcgen")] fn dtls12_handles_duplicate_packets() { diff --git a/tests/dtls12/retransmit.rs b/tests/dtls12/retransmit.rs index b1c9dde8..eff9df01 100644 --- a/tests/dtls12/retransmit.rs +++ b/tests/dtls12/retransmit.rs @@ -122,6 +122,10 @@ fn prepare_server_final_flight_resend(max_queue_rx: usize) -> FinalFlightResend "resend flight 6 should include epoch 1 Finished" ); assert_epochs_and_seq_increased(&f6_init_hdrs, &f6_resend_hdrs); + server + .handle_timeout(now) + .expect("arm server resend cooldown"); + assert!(collect_packets(&mut server).is_empty()); FinalFlightResend { client, @@ -467,6 +471,12 @@ fn dtls12_final_flight_resend_can_share_post_release_appdata_tail() { .. } = prepare_server_final_flight_resend(Config::default().max_queue_rx()); + let cooldown = drain_outputs(&mut server).timeout.expect("resend cooldown"); + server + .handle_timeout(cooldown) + .expect("expire previous resend cooldown"); + assert!(collect_packets(&mut server).is_empty()); + let filler = vec![0x55; 50]; for i in 0..9 { server @@ -1101,6 +1111,12 @@ fn dtls12_stale_client_hello_before_peer_confirmed_triggers_resend() { .. } = prepare_server_final_flight_resend(RX_QUEUE_LIMIT); + let cooldown = drain_outputs(&mut server).timeout.expect("resend cooldown"); + server + .handle_timeout(cooldown) + .expect("expire previous resend cooldown"); + assert!(collect_packets(&mut server).is_empty()); + // The server has sent flight 6 but has no confirmation the client received // it (no authenticated application data has arrived). A stale ClientHello // here means "peer still needs my flight" and must drive a resend. From bdcf7baf23022d9b492412a562c58490876fb717 Mon Sep 17 00:00:00 2001 From: Martin Algesten Date: Fri, 2 Oct 2026 10:12:54 +0200 Subject: [PATCH 2/4] Document DTLS 1.2 datagram validation and resend backoff --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 65cb6928..c67130dc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,6 @@ # Unreleased + * Reject internally reordered DTLS 1.2 datagrams and back off duplicate flight resends #169 * Restore DTLS buffer reuse and reject unpooled buffer returns #167 # 0.7.4 From e2e6c491d8f81486eb6b3a1d11e412490dfd4fbc Mon Sep 17 00:00:00 2001 From: Martin Algesten Date: Fri, 2 Oct 2026 10:32:36 +0200 Subject: [PATCH 3/4] Reject unordered datagrams with an internal parse error --- src/dtls12/incoming.rs | 41 ++++++++++++++++++----------------------- 1 file changed, 18 insertions(+), 23 deletions(-) diff --git a/src/dtls12/incoming.rs b/src/dtls12/incoming.rs index e2a7e81d..135d4ee2 100644 --- a/src/dtls12/incoming.rs +++ b/src/dtls12/incoming.rs @@ -3,6 +3,7 @@ use std::ops::Deref; use std::sync::atomic::{AtomicBool, Ordering}; use arrayvec::ArrayVec; +use nom::error::ErrorKind; use crate::buffer::{Buf, TmpBuf}; use crate::crypto::{Aad, Nonce}; @@ -80,18 +81,6 @@ impl Records { return Err(error); } - // The receive queue preserves each datagram's internal order. Discard an - // unordered datagram before classifying controls or retaining any of it, - // so a subsequent clean retransmission can recover. - if !Self::is_ordered(&parsed_records, decrypt.is_peer_encryption_enabled()) { - for record in parsed_records { - decrypt.push_buffer(record.into_buffer()); - } - return Ok(Records { - records: ArrayVec::new(), - }); - } - let mut records = ArrayVec::new(); for record in parsed_records { if let Some(record) = decrypt.classify_record(record)? { @@ -139,6 +128,14 @@ impl Records { packet = &packet[record_end..]; } + // The receive queue preserves each datagram's internal order. Reject an + // unordered datagram before classifying controls or retaining any of it, + // so a subsequent clean retransmission can recover. The caller recycles + // all parsed records through the normal parser-error cleanup path. + if !Self::is_ordered(parsed_records, decrypt.is_peer_encryption_enabled()) { + return Err(InternalError::parse(ErrorKind::Verify)); + } + Ok(()) } @@ -601,9 +598,9 @@ mod tests { } let mut handler = TestHandler::default(); - let incoming = Incoming::parse_packet(&packet, &mut handler, None) - .expect("unordered datagrams are silently discarded"); - assert!(incoming.is_none()); + let error = Incoming::parse_packet(&packet, &mut handler, None) + .expect_err("unordered datagrams produce an internal parse error"); + assert!(error.into_public_error().is_none()); assert_eq!(handler.classify_calls, 0); assert_eq!(handler.buffers_returned, handler.buffers_acquired); } @@ -618,11 +615,9 @@ mod tests { handshakes.extend_from_slice(&build_handshake(0, 0, &[0; 4])); let packet = build_record(ContentType::Handshake, 0, 0, &handshakes); let mut handler = TestHandler::default(); - assert!( - Incoming::parse_packet(&packet, &mut handler, None) - .expect("unordered datagrams are silently discarded") - .is_none() - ); + let error = Incoming::parse_packet(&packet, &mut handler, None) + .expect_err("unordered datagrams produce an internal parse error"); + assert!(error.into_public_error().is_none()); assert_eq!(handler.classify_calls, 0); assert_eq!(handler.buffers_returned, handler.buffers_acquired); } @@ -661,9 +656,9 @@ mod tests { ); for packet in [[finished, ccs.clone()].concat(), [ccs, handshake].concat()] { let mut handler = TestHandler::default(); - let incoming = Incoming::parse_packet(&packet, &mut handler, None) - .expect("unordered datagrams are silently discarded"); - assert!(incoming.is_none()); + let error = Incoming::parse_packet(&packet, &mut handler, None) + .expect_err("unordered datagrams produce an internal parse error"); + assert!(error.into_public_error().is_none()); assert_eq!(handler.classify_calls, 0); assert_eq!(handler.buffers_returned, handler.buffers_acquired); } From 768b0bc4025cbf2547f07150e51683714109261c Mon Sep 17 00:00:00 2001 From: Martin Algesten Date: Fri, 2 Oct 2026 10:37:32 +0200 Subject: [PATCH 4/4] Split the DTLS 1.2 changelog entry --- CHANGELOG.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c67130dc..a2e65543 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,7 @@ # Unreleased - * Reject internally reordered DTLS 1.2 datagrams and back off duplicate flight resends #169 + * Reject internally reordered DTLS 1.2 datagrams #169 + * Back off duplicate DTLS 1.2 flight resends #169 * Restore DTLS buffer reuse and reject unpooled buffer returns #167 # 0.7.4