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
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
# Unreleased

* 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
Expand Down
232 changes: 207 additions & 25 deletions src/dtls12/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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 {
Expand All @@ -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));
}
Expand Down Expand Up @@ -525,38 +541,33 @@ 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) {
debug!("Begin flight {}", flight_no);
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) {
Expand All @@ -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;
}
}

Expand All @@ -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);

Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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;
}
}

Expand Down Expand Up @@ -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<u8> {
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);
Expand Down
Loading
Loading