From 780bd5bdb38e702cca841026467c5cb28aa49743 Mon Sep 17 00:00:00 2001 From: Andrey Mnatsakanov Date: Mon, 21 Sep 2026 19:04:48 +0200 Subject: [PATCH 1/3] feat(turn): request the relay transport independently of the server transport MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Allocate derived REQUESTED-TRANSPORT from the transport used to reach the server, so a client talking to its server over TCP asked for an RFC 6062 TCP allocation. RFC 8656 §7.1 keeps the two apart: REQUESTED-TRANSPORT names what the allocation relays, and a client on TCP normally still wants a UDP relay. ClientConfig gains `requested_transport`, defaulting to UDP, used by the Allocate and by its authenticated retry. Callers that set nothing keep what they had; callers that build ClientConfig as a struct literal need the new field, which the in-tree examples and tests now set. Refs webrtc-rs/webrtc#848. --- .../trickle-ice-relay/trickle-ice-relay.rs | 1 + examples/trickle-ice/trickle-ice.rs | 1 + rtc-turn/examples/turn_client_udp.rs | 1 + rtc-turn/src/client/client_test.rs | 51 +++++++++++++++++++ rtc-turn/src/client/mod.rs | 34 +++++++++---- rtc-turn/src/lib.rs | 3 +- 6 files changed, 79 insertions(+), 12 deletions(-) diff --git a/examples/trickle-ice-relay/trickle-ice-relay.rs b/examples/trickle-ice-relay/trickle-ice-relay.rs index 4856a467..e20ad82c 100644 --- a/examples/trickle-ice-relay/trickle-ice-relay.rs +++ b/examples/trickle-ice-relay/trickle-ice-relay.rs @@ -183,6 +183,7 @@ async fn run_main_loop( turn_serv_addr: turn_server_addr, local_addr, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: cred[0].to_string(), password: cred[1].to_string(), realm: turn_realm.to_string(), diff --git a/examples/trickle-ice/trickle-ice.rs b/examples/trickle-ice/trickle-ice.rs index 6fe5a173..4bd7cee0 100644 --- a/examples/trickle-ice/trickle-ice.rs +++ b/examples/trickle-ice/trickle-ice.rs @@ -280,6 +280,7 @@ async fn run_main_loop(cli: Cli) -> Result<()> { turn_serv_addr: turn_server_str.clone(), local_addr, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: cred[0].to_string(), password: cred[1].to_string(), realm: cli.turn_realm.to_string(), diff --git a/rtc-turn/examples/turn_client_udp.rs b/rtc-turn/examples/turn_client_udp.rs index 5d502a8e..fd373ddc 100644 --- a/rtc-turn/examples/turn_client_udp.rs +++ b/rtc-turn/examples/turn_client_udp.rs @@ -87,6 +87,7 @@ fn main() -> Result<()> { turn_serv_addr: turn_server_addr, local_addr, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: cred[0].to_string(), password: cred[1].to_string(), realm: realm.to_string(), diff --git a/rtc-turn/src/client/client_test.rs b/rtc-turn/src/client/client_test.rs index d16dc9d1..898bfda2 100644 --- a/rtc-turn/src/client/client_test.rs +++ b/rtc-turn/src/client/client_test.rs @@ -17,6 +17,7 @@ fn create_listening_test_client(rto_in_ms: u64) -> Result<(UdpSocket, Client)> { turn_serv_addr: String::new(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: String::new(), password: String::new(), realm: String::new(), @@ -39,6 +40,7 @@ fn create_listening_test_client_with_stun_serv() -> Result<(UdpSocket, Client)> turn_serv_addr: String::new(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: String::new(), password: String::new(), realm: String::new(), @@ -216,6 +218,7 @@ fn test_relay_refresh_timers_run_on_injected_time() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -297,6 +300,7 @@ fn test_allocation_refresh_interval_cap_is_applied() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -360,6 +364,7 @@ fn test_overdue_relay_refreshes_are_rescheduled_from_now() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -427,6 +432,7 @@ fn test_zero_lifetime_relay_does_not_freeze_its_deadline() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -486,6 +492,7 @@ fn test_refresh_response_reschedules_shorter_lifetime_from_response_time() -> Re turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -538,6 +545,7 @@ fn test_refresh_response_reschedule_applies_interval_cap() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -596,6 +604,7 @@ fn test_zero_lifetime_response_drops_the_relay() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -664,6 +673,7 @@ fn test_zero_lifetime_allocate_response_is_an_error() -> Result<()> { turn_serv_addr: "127.0.0.1:3478".to_owned(), local_addr: udp_socket.local_addr()?, transport_protocol: TransportProtocol::UDP, + requested_transport: TransportProtocol::UDP, username: "user".to_owned(), password: "pass".to_owned(), realm: "realm".to_owned(), @@ -704,3 +714,44 @@ fn test_zero_lifetime_allocate_response_is_an_error() -> Result<()> { } } } + +/// Over a TCP connection to the server the relay is still UDP: `REQUESTED-TRANSPORT` names what +/// the allocation relays, not how the client reaches the server (RFC 8656 §7.1). Deriving it from +/// the control transport asks a TCP client's server for an RFC 6062 TCP allocation instead. +#[test] +fn test_allocate_over_tcp_requests_a_udp_relay() -> Result<()> { + let mut client = Client::new( + ClientConfig { + turn_serv_addr: "127.0.0.1:3478".to_owned(), + local_addr: "127.0.0.1:50000".parse().unwrap(), + transport_protocol: TransportProtocol::TCP, + requested_transport: TransportProtocol::UDP, + ..Default::default() + }, + test_crypto_provider(), + )?; + client.allocate(Instant::now())?; + + let transmit = client.poll_write().expect("an Allocate request"); + assert_eq!( + transmit.transport.transport_protocol, + TransportProtocol::TCP, + "the request itself still travels over TCP" + ); + let mut msg = Message::new(); + msg.raw = transmit.message.to_vec(); + msg.decode()?; + let mut requested = RequestedTransport::default(); + requested.get_from(&msg)?; + assert_eq!(requested.protocol, PROTO_UDP, "the relay is UDP"); + Ok(()) +} + +/// Existing callers that set nothing keep asking for what they always got. +#[test] +fn test_requested_transport_defaults_to_udp() { + assert_eq!( + ClientConfig::default().requested_transport, + TransportProtocol::UDP + ); +} diff --git a/rtc-turn/src/client/mod.rs b/rtc-turn/src/client/mod.rs index 5ea2a33c..fbf15d33 100644 --- a/rtc-turn/src/client/mod.rs +++ b/rtc-turn/src/client/mod.rs @@ -45,7 +45,7 @@ use crate::proto::lifetime::Lifetime; use crate::proto::peeraddr::*; use crate::proto::relayaddr::RelayedAddress; use crate::proto::reqtrans::RequestedTransport; -use crate::proto::{PROTO_TCP, PROTO_UDP}; +use crate::proto::{PROTO_TCP, PROTO_UDP, Protocol}; use shared::error::{Error, Result}; use shared::util::lookup_host; use shared::{TransportContext, TransportMessage, TransportProtocol}; @@ -122,6 +122,12 @@ pub struct ClientConfig { pub local_addr: SocketAddr, /// Whether to reach the server over UDP or TCP. pub transport_protocol: TransportProtocol, + /// The transport the allocation relays, sent as `REQUESTED-TRANSPORT`. + /// + /// Independent of `transport_protocol`: a client that reaches the server over TCP still + /// wants a UDP relay (RFC 8656 §7.1). Asking for TCP here requests an RFC 6062 TCP + /// allocation instead, which this client does not implement. + pub requested_transport: TransportProtocol, /// The long-term credential username for the TURN server. pub username: String, /// The long-term credential password. @@ -147,6 +153,7 @@ impl Default for ClientConfig { turn_serv_addr: "".to_string(), local_addr: SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, 0)), transport_protocol: Default::default(), + requested_transport: Default::default(), username: "".to_string(), password: "".to_string(), realm: "".to_string(), @@ -164,6 +171,7 @@ pub struct Client { turn_serv_addr: Option, local_addr: SocketAddr, transport_protocol: TransportProtocol, + requested_transport: TransportProtocol, username: Username, password: String, realm: Realm, @@ -207,6 +215,7 @@ impl Client { turn_serv_addr, local_addr: config.local_addr, transport_protocol: config.transport_protocol, + requested_transport: config.requested_transport, username: Username::new(ATTR_USERNAME, config.username), password: config.password, realm: Realm::new(ATTR_REALM, config.realm), @@ -589,6 +598,17 @@ impl Client { Ok(()) } + /// The `REQUESTED-TRANSPORT` value: what the allocation relays, which is not necessarily + /// how this client reaches the server. + fn requested_protocol(&self) -> Protocol { + match self.requested_transport { + TransportProtocol::TCP => PROTO_TCP, + // UDP, and anything a later transport adds: a UDP relay is what TURN allocates + // unless a TCP one is asked for by name. + _ => PROTO_UDP, + } + } + /// Allocate sends a TURN allocation request to the given transport address pub fn allocate(&mut self, now: Instant) -> Result { let mut msg = Message::new(); @@ -596,11 +616,7 @@ impl Client { Box::new(TransactionId::new()), Box::new(MessageType::new(METHOD_ALLOCATE, CLASS_REQUEST)), Box::new(RequestedTransport { - protocol: if self.transport_protocol == TransportProtocol::UDP { - PROTO_UDP - } else { - PROTO_TCP - }, + protocol: self.requested_protocol(), }), Box::new(FINGERPRINT), ])?; @@ -661,11 +677,7 @@ impl Client { Box::new(tid), Box::new(MessageType::new(METHOD_ALLOCATE, CLASS_REQUEST)), Box::new(RequestedTransport { - protocol: if self.transport_protocol == TransportProtocol::UDP { - PROTO_UDP - } else { - PROTO_TCP - }, + protocol: self.requested_protocol(), }), Box::new(self.username.clone()), Box::new(self.realm.clone()), diff --git a/rtc-turn/src/lib.rs b/rtc-turn/src/lib.rs index f6714c15..7388bd10 100644 --- a/rtc-turn/src/lib.rs +++ b/rtc-turn/src/lib.rs @@ -29,7 +29,8 @@ //! let config = ClientConfig { //! turn_serv_addr: "turn.example.com:3478".to_owned(), //! local_addr: "0.0.0.0:0".parse().unwrap(), -//! transport_protocol: TransportProtocol::UDP, +//! transport_protocol: TransportProtocol::UDP, // how the server is reached +//! requested_transport: TransportProtocol::UDP, // what the allocation relays //! username: "user".to_owned(), //! password: "pass".to_owned(), //! realm: "example.com".to_owned(), From 1454debb60bb86f97f1ac85b75130af71ae41498 Mon Sep 17 00:00:00 2001 From: Andrey Mnatsakanov Date: Mon, 21 Sep 2026 19:06:58 +0200 Subject: [PATCH 2/3] feat(shared): split a TURN byte stream into STUN and ChannelData messages MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Over TCP, TURN sends its messages back to back with no framing of its own (RFC 8656 §12.5): a STUN message is 20 bytes plus its length field, a ChannelData message 4 bytes plus its length rounded up to a multiple of 4. tcp_framing is RFC 4571's length prefix, which is ICE-TCP's framing and not TURN's, so a host reading a TURN stream had nothing to split it with. TurnStreamDecoder has TcpFrameDecoder's shape — push bytes, pop whole messages — so a host can hold either per stream. Messages come out as they were on the wire, padding included; a first byte that starts neither kind is an error, since a stream has nowhere to resynchronise. Refs webrtc-rs/webrtc#848. --- rtc-shared/src/lib.rs | 4 + rtc-shared/src/turn_framing.rs | 187 +++++++++++++++++++++++++++++++++ 2 files changed, 191 insertions(+) create mode 100644 rtc-shared/src/turn_framing.rs diff --git a/rtc-shared/src/lib.rs b/rtc-shared/src/lib.rs index c6c5507f..190f7e06 100644 --- a/rtc-shared/src/lib.rs +++ b/rtc-shared/src/lib.rs @@ -22,6 +22,8 @@ //! * [`replay_detector`] — replay protection shared by DTLS and SRTP. (Cryptography //! itself lives in the separate `rtc-crypto` crate, behind `RTCCryptoProvider`.) //! * [`tcp_framing`] — RFC 4571 length-prefixed framing, for ICE-TCP candidates. +//! * [`turn_framing`] — splitting a TURN stream into its self-delimiting messages, for TURN +//! over TCP. //! * [`ifaces`] — local interface enumeration used during ICE candidate gathering. //! //! # Feature flags @@ -78,6 +80,8 @@ pub mod tcp_framing; /// Conversions between monotonic, Unix and NTP time. pub mod time; pub(crate) mod transport; +/// Splitting a TURN-over-TCP byte stream into STUN and ChannelData messages. +pub mod turn_framing; /// Small shared helpers: packet demultiplexing predicates and random-string generation. pub mod util; diff --git a/rtc-shared/src/turn_framing.rs b/rtc-shared/src/turn_framing.rs new file mode 100644 index 00000000..9148d09b --- /dev/null +++ b/rtc-shared/src/turn_framing.rs @@ -0,0 +1,187 @@ +//! Splitting a TURN byte stream back into messages (RFC 8656 §12.5). +//! +//! Over TCP or TLS, TURN sends STUN messages and ChannelData messages back to +//! back with no framing of their own: each is self-delimiting. This is not the +//! RFC 4571 length prefix that ICE-TCP uses — see [`tcp_framing`](crate::tcp_framing) +//! for that — and a stream carries one or the other, never both. +//! +//! ```text +//! STUN 00xxxxxx ........ | length (16) | magic cookie, id ... | attributes +//! total = 20 + length +//! ChannelData 01xxxxxx ........ | length (16) | data | padding to a multiple of 4 +//! total = 4 + length rounded up to 4 +//! ``` +//! +//! The first two bits tell them apart. Anything else cannot be TURN, and since +//! a stream has no resynchronisation point it is an error that ends the stream. +//! +//! ```rust +//! use rtc_shared::turn_framing::TurnStreamDecoder; +//! +//! // A ChannelData message on channel 0x4000 carrying three bytes, padded to 8. +//! let wire = [0x40, 0x00, 0x00, 0x03, b'a', b'b', b'c', 0x00]; +//! let mut decoder = TurnStreamDecoder::new(); +//! decoder.extend_from_slice(&wire[..5]); +//! assert_eq!(decoder.next_message().unwrap(), None); // not all of it yet +//! decoder.extend_from_slice(&wire[5..]); +//! assert_eq!(decoder.next_message().unwrap(), Some(wire.to_vec())); +//! ``` + +use crate::error::{Error, Result}; + +/// Sans-IO splitter for a TURN stream: push bytes as they arrive, pop whole +/// messages. A message comes out exactly as it was on the wire, ChannelData +/// padding included, which TURN's decoders accept. +#[derive(Debug, Default)] +pub struct TurnStreamDecoder { + buffer: Vec, +} + +impl TurnStreamDecoder { + /// Creates a decoder with an empty buffer. + pub fn new() -> Self { + Self::default() + } + + /// Appends bytes read from the stream. + pub fn extend_from_slice(&mut self, data: &[u8]) { + self.buffer.extend_from_slice(data); + } + + /// The next complete message, `Ok(None)` while one is still arriving, or an + /// error when the stream holds something that is not TURN. + pub fn next_message(&mut self) -> Result>> { + // Both kinds carry their length in bytes 2..4, so four bytes decide it. + if self.buffer.len() < CHANNEL_DATA_HEADER_LEN { + return Ok(None); + } + let length = u16::from_be_bytes([self.buffer[2], self.buffer[3]]) as usize; + let total = match self.buffer[0] >> 6 { + 0b00 => STUN_HEADER_LEN + length, + 0b01 => CHANNEL_DATA_HEADER_LEN + length.next_multiple_of(4), + _ => { + return Err(Error::OtherTurnErr(format!( + "stream byte {:#04x} starts neither a STUN nor a ChannelData message", + self.buffer[0] + ))); + } + }; + if self.buffer.len() < total { + return Ok(None); + } + let message = self.buffer[..total].to_vec(); + self.buffer.drain(..total); + Ok(Some(message)) + } + + /// Bytes held that do not yet make a whole message. + pub fn buffered_len(&self) -> usize { + self.buffer.len() + } +} + +/// A STUN header: type, length, magic cookie and transaction id (RFC 8489 §5). +const STUN_HEADER_LEN: usize = 20; +/// A ChannelData header: channel number and length (RFC 8656 §12.4). +const CHANNEL_DATA_HEADER_LEN: usize = 4; + +#[cfg(test)] +mod tests { + use super::*; + + /// A STUN Binding request header with `length` bytes of attributes (zeros). + fn stun(length: u16) -> Vec { + let mut msg = vec![0x00, 0x01]; + msg.extend_from_slice(&length.to_be_bytes()); + msg.extend_from_slice(&[0x21, 0x12, 0xa4, 0x42]); + msg.extend_from_slice(&[7; 12]); + msg.extend(std::iter::repeat_n(0u8, length as usize)); + msg + } + + /// A ChannelData message on 0x4000 carrying `data`, padded to a multiple of 4. + fn channel_data(data: &[u8]) -> Vec { + let mut msg = vec![0x40, 0x00]; + msg.extend_from_slice(&(data.len() as u16).to_be_bytes()); + msg.extend_from_slice(data); + while msg.len() % 4 != 0 { + msg.push(0); + } + msg + } + + fn drain(decoder: &mut TurnStreamDecoder) -> Vec> { + let mut out = Vec::new(); + while let Some(msg) = decoder.next_message().unwrap() { + out.push(msg); + } + out + } + + #[test] + fn a_stun_message_is_twenty_bytes_plus_its_length() { + let msg = stun(8); + let mut decoder = TurnStreamDecoder::new(); + decoder.extend_from_slice(&msg); + assert_eq!(drain(&mut decoder), vec![msg]); + assert_eq!(decoder.buffered_len(), 0); + } + + #[test] + fn channel_data_comes_out_with_its_padding() { + for data in [&b""[..], b"a", b"ab", b"abc", b"abcd", b"abcde"] { + let msg = channel_data(data); + let mut decoder = TurnStreamDecoder::new(); + decoder.extend_from_slice(&msg); + assert_eq!(drain(&mut decoder), vec![msg], "{} data bytes", data.len()); + assert_eq!(decoder.buffered_len(), 0); + } + } + + #[test] + fn messages_arriving_together_come_out_one_by_one() { + let (a, b, c) = (stun(4), channel_data(b"xyz"), stun(0)); + let mut decoder = TurnStreamDecoder::new(); + decoder.extend_from_slice(&[a.clone(), b.clone(), c.clone()].concat()); + assert_eq!(drain(&mut decoder), vec![a, b, c]); + } + + #[test] + fn a_message_split_across_reads_waits_until_it_is_whole() { + let msg = channel_data(b"hello"); + let mut decoder = TurnStreamDecoder::new(); + // One byte at a time, including splits inside the four-byte header. + for (i, byte) in msg.iter().enumerate() { + decoder.extend_from_slice(&[*byte]); + if i + 1 < msg.len() { + assert_eq!(decoder.next_message().unwrap(), None, "early at byte {i}"); + } + } + assert_eq!(drain(&mut decoder), vec![msg]); + } + + #[test] + fn a_partial_message_after_a_whole_one_stays_buffered() { + let whole = stun(0); + let next = channel_data(b"later"); + let mut decoder = TurnStreamDecoder::new(); + decoder.extend_from_slice(&[whole.clone(), next[..3].to_vec()].concat()); + assert_eq!(drain(&mut decoder), vec![whole]); + assert_eq!(decoder.buffered_len(), 3); + decoder.extend_from_slice(&next[3..]); + assert_eq!(drain(&mut decoder), vec![next]); + } + + #[test] + fn bytes_that_are_not_turn_end_the_stream() { + // 10xxxxxx and 11xxxxxx are neither STUN nor ChannelData. + for first in [0x80u8, 0xbf, 0xc0, 0xff] { + let mut decoder = TurnStreamDecoder::new(); + decoder.extend_from_slice(&[first, 0, 0, 4, 1, 2, 3, 4]); + assert!( + matches!(decoder.next_message(), Err(Error::OtherTurnErr(_))), + "first byte {first:#04x} was accepted" + ); + } + } +} From 5d14a40fd4b9d68aa40dbf776877f02a951febd6 Mon Sep 17 00:00:00 2001 From: Andrey Mnatsakanov Date: Mon, 21 Sep 2026 19:23:29 +0200 Subject: [PATCH 3/3] fix(turn): do not retransmit requests over TCP MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Transactions retransmitted on every transport: 200 ms doubling to 1.6 s, seven requests, a timeout after about eight seconds. Over a reliable transport RFC 8489 §6.2.2 has the client send once and wait Ti, 39.5 s — TCP is already retransmitting underneath, and repeating the request on top only sends the server duplicates of it. A transaction over anything that is not UDP now has one deadline at Ti and no retransmission; passing it reports TransactionTimeout as before. UDP is unchanged. Kept as its own commit: it is not needed for TURN over TCP to work, only for it to behave as the RFC asks. Refs webrtc-rs/webrtc#848. --- rtc-turn/src/client/client_test.rs | 77 ++++++++++++++++++++++++++++++ rtc-turn/src/client/transaction.rs | 22 ++++++++- 2 files changed, 98 insertions(+), 1 deletion(-) diff --git a/rtc-turn/src/client/client_test.rs b/rtc-turn/src/client/client_test.rs index 898bfda2..08cf0abd 100644 --- a/rtc-turn/src/client/client_test.rs +++ b/rtc-turn/src/client/client_test.rs @@ -755,3 +755,80 @@ fn test_requested_transport_defaults_to_udp() { TransportProtocol::UDP ); } + +/// Over a reliable transport a request is sent once and answered or timed out as a whole: TCP +/// already retransmits, and repeating the request on top of it only sends the server duplicates +/// (RFC 8489 §6.2.2). The wait is Ti, 39.5 seconds, not the UDP schedule's ~8. +#[test] +fn test_a_request_over_tcp_is_sent_once_and_times_out_at_ti() -> Result<()> { + let start = Instant::now(); + let mut client = Client::new( + ClientConfig { + turn_serv_addr: "127.0.0.1:3478".to_owned(), + local_addr: "127.0.0.1:50000".parse().unwrap(), + transport_protocol: TransportProtocol::TCP, + ..Default::default() + }, + test_crypto_provider(), + )?; + client.allocate(start)?; + // allocate() returns the id its authenticated retry will use; the timeout names the + // request that was actually sent, so take the id from the wire. + let sent = client.poll_write().expect("the request goes out once"); + let mut request = Message::new(); + request.raw = sent.message.to_vec(); + request.decode()?; + let tid = request.transaction_id; + + // Walk every deadline the client asks for, up to just short of Ti. + let ti = Duration::from_millis(39_500); + let mut sends = 0; + while let Some(deadline) = client.poll_timeout() { + if deadline >= start + ti { + break; + } + client.handle_timeout(deadline)?; + while client.poll_write().is_some() { + sends += 1; + } + } + assert_eq!(sends, 0, "nothing is retransmitted over TCP"); + assert!( + client.poll_event().is_none(), + "no timeout before Ti has passed" + ); + + client.handle_timeout(start + ti)?; + assert!(client.poll_write().is_none(), "not even at the deadline"); + assert!( + matches!(client.poll_event(), Some(Event::TransactionTimeout(id)) if id == tid), + "at Ti the request has timed out" + ); + Ok(()) +} + +/// UDP keeps its schedule: the first retransmission is due one RTO after the request. +#[test] +fn test_a_request_over_udp_is_still_retransmitted() -> Result<()> { + let start = Instant::now(); + let mut client = Client::new( + ClientConfig { + turn_serv_addr: "127.0.0.1:3478".to_owned(), + local_addr: "127.0.0.1:50000".parse().unwrap(), + transport_protocol: TransportProtocol::UDP, + ..Default::default() + }, + test_crypto_provider(), + )?; + client.allocate(start)?; + assert!(client.poll_write().is_some()); + let first = client.poll_timeout().expect("a retransmit deadline"); + assert!( + first < start + Duration::from_secs(1), + "{:?}", + first - start + ); + client.handle_timeout(first)?; + assert!(client.poll_write().is_some(), "UDP retransmits"); + Ok(()) +} diff --git a/rtc-turn/src/client/transaction.rs b/rtc-turn/src/client/transaction.rs index d02d29cf..74049e84 100644 --- a/rtc-turn/src/client/transaction.rs +++ b/rtc-turn/src/client/transaction.rs @@ -13,6 +13,9 @@ use stun::textattrs::TextAttribute; const MAX_RTX_INTERVAL_IN_MS: u64 = 1600; const MAX_RTX_COUNT: u16 = 7; // total 7 requests (Rc) +/// Ti (RFC 8489 §6.2.2): over a reliable transport a request is sent once, and +/// this is how long the client waits for its answer. +const RELIABLE_TRANSACTION_TIMEOUT: Duration = Duration::from_millis(39_500); pub(crate) enum TransactionType { BindingRequest, @@ -62,7 +65,11 @@ impl Transaction { transport_protocol: config.transport_protocol, n_rtx: 0, interval: config.interval, - timeout: config.now.add(Duration::from_millis(config.interval)), + timeout: if is_reliable(config.transport_protocol) { + config.now.add(RELIABLE_TRANSACTION_TIMEOUT) + } else { + config.now.add(Duration::from_millis(config.interval)) + }, transmits: VecDeque::new(), } } @@ -77,6 +84,13 @@ impl Transaction { pub(crate) fn handle_timeout(&mut self, now: Instant) { if self.retries() < MAX_RTX_COUNT && self.timeout <= now { + if is_reliable(self.transport_protocol) { + // One deadline and no retransmission: TCP already retransmits, and + // repeating the request over it only sends the server duplicates. + // Passing Ti exhausts the transaction, which reports the timeout. + self.n_rtx = MAX_RTX_COUNT; + return; + } self.n_rtx += 1; self.interval *= 2; if self.interval > MAX_RTX_INTERVAL_IN_MS { @@ -121,6 +135,12 @@ impl Transaction { } } +/// Whether the transport delivers reliably on its own, so a request must not be +/// retransmitted over it. Anything that is not UDP is a stream. +fn is_reliable(transport_protocol: TransportProtocol) -> bool { + transport_protocol != TransportProtocol::UDP +} + // TransactionMap is a thread-safe transaction map #[derive(Default)] pub(crate) struct TransactionMap {