Skip to main content

timestamped_socket/socket/
linux.rs

1use std::{
2    marker::PhantomData,
3    net::{Ipv4Addr, Ipv6Addr, SocketAddr, SocketAddrV4, SocketAddrV6},
4    sync::Arc,
5};
6
7use tokio::io::{unix::AsyncFd, Interest};
8
9use crate::{
10    control_message::{ControlMessage, MessageQueue, EXPECTED_MAX_CMSG_SIZE},
11    interface::{lookup_phc, InterfaceName},
12    networkaddress::{sealed::PrivateToken, EthernetAddress, MacAddress, NetworkAddress},
13    raw_socket::RawSocket,
14    socket::TimestampData,
15};
16
17use super::{InterfaceTimestampMode, Open, Socket};
18
19const SOF_TIMESTAMPING_BIND_PHC: libc::c_uint = 1 << 15;
20
21impl<A: NetworkAddress, S> Socket<A, S> {
22    pub(super) async fn fetch_send_timestamp(
23        socket: Arc<AsyncFd<RawSocket>>,
24    ) -> std::io::Result<(u32, TimestampData)> {
25        let try_read = |socket: &RawSocket| fetch_send_timestamp_try_read(socket);
26
27        loop {
28            // the timestamp being available triggers the error interest
29            match socket.async_io(Interest::ERROR, try_read).await? {
30                Some((counter, timestamp_data)) => break Ok((counter, timestamp_data)),
31                None => continue,
32            }
33        }
34    }
35}
36
37/// This function tries to fetch a send timestamp from a single error queue message.
38///
39/// We assume that we get error queue messages for exactly two reasons:
40///  - When the driver has pushed a timestamp set up, in which timestamps will be present
41///  - When something has gone wrong in the send process after the return of the send*
42///    family of functions (such as an ICMP error).
43///
44/// In particular, we assume that the second scenario occurs independent of whether the first
45/// occurs, and therefore we fully ignore those messages, only logging them but never
46/// returning the message index, to avoid upper layers of interpreting the error as a marker
47/// that the timestamp is definitively unavailable.
48///
49/// Note that this means that we don't expect, and therefore don't report, any signal from
50/// the kernel that the message has been sent but there won't be timestamps available. This
51/// scenario is handled in the higher layers through timeouts.
52///
53/// The above are assumptions for us as the kernel does not clearly document precisely how
54/// this works, and we haven't had the capacity to do a full deep dive into the kernel source.
55fn fetch_send_timestamp_try_read(
56    socket: &RawSocket,
57) -> std::io::Result<Option<(u32, TimestampData)>> {
58    let mut control_buf = [0; EXPECTED_MAX_CMSG_SIZE];
59
60    // NOTE: this read could block!
61    let (_, control_messages, _) =
62        socket.receive_message(&mut [], &mut control_buf, MessageQueue::Error)?;
63
64    let mut send_ts = None;
65    let mut counter = None;
66    for msg in control_messages {
67        match msg {
68            ControlMessage::Timestamping { software, hardware } => {
69                send_ts = Some(TimestampData {
70                    software,
71                    hardware,
72                    timestamp_mode: InterfaceTimestampMode::default(),
73                });
74            }
75
76            ControlMessage::ReceiveError(error) => {
77                // the timestamping does not set a message; if there is a message, that means
78                // something else is wrong, and we want to know about it.
79                if error.ee_errno as libc::c_int != libc::ENOMSG {
80                    tracing::debug!(error.ee_data, "error message on the MSG_ERRQUEUE");
81                }
82
83                counter = Some(error.ee_data);
84            }
85
86            ControlMessage::DestinationIp(_) => {
87                tracing::debug!("unexpected destination ip control message");
88            }
89
90            ControlMessage::Other(msg) => {
91                tracing::debug!(
92                    msg.cmsg_level,
93                    msg.cmsg_type,
94                    "unexpected message on the MSG_ERRQUEUE",
95                );
96            }
97        }
98    }
99
100    Ok(counter.zip(send_ts))
101}
102
103pub(super) fn configure_timestamping(
104    socket: &RawSocket,
105    interface: Option<InterfaceName>,
106    mode: InterfaceTimestampMode,
107    mut bind_phc: Option<u32>,
108) -> std::io::Result<()> {
109    // Check if the phc is not the interface-native phc.
110    if let Some(interface) = interface {
111        if lookup_phc(interface) == bind_phc {
112            bind_phc = None
113        }
114    }
115
116    let options = match mode {
117        InterfaceTimestampMode::HardwareAll | InterfaceTimestampMode::HardwarePTPAll => {
118            libc::SOF_TIMESTAMPING_RAW_HARDWARE
119                | libc::SOF_TIMESTAMPING_RX_SOFTWARE
120                | libc::SOF_TIMESTAMPING_TX_SOFTWARE
121                | libc::SOF_TIMESTAMPING_RX_HARDWARE
122                | libc::SOF_TIMESTAMPING_TX_HARDWARE
123                | libc::SOF_TIMESTAMPING_OPT_TSONLY
124                | libc::SOF_TIMESTAMPING_OPT_ID
125                | bind_phc
126                    .map(|_| SOF_TIMESTAMPING_BIND_PHC)
127                    .unwrap_or_default()
128        }
129        InterfaceTimestampMode::HardwareRecv | InterfaceTimestampMode::HardwarePTPRecv => {
130            libc::SOF_TIMESTAMPING_RAW_HARDWARE
131                | libc::SOF_TIMESTAMPING_RX_SOFTWARE
132                | libc::SOF_TIMESTAMPING_RX_HARDWARE
133                | bind_phc
134                    .map(|_| SOF_TIMESTAMPING_BIND_PHC)
135                    .unwrap_or_default()
136        }
137        InterfaceTimestampMode::SoftwareAll => {
138            libc::SOF_TIMESTAMPING_SOFTWARE
139                | libc::SOF_TIMESTAMPING_RX_SOFTWARE
140                | libc::SOF_TIMESTAMPING_TX_SOFTWARE
141                | libc::SOF_TIMESTAMPING_OPT_TSONLY
142                | libc::SOF_TIMESTAMPING_OPT_ID
143        }
144        InterfaceTimestampMode::SoftwareRecv => {
145            libc::SOF_TIMESTAMPING_SOFTWARE | libc::SOF_TIMESTAMPING_RX_SOFTWARE
146        }
147        InterfaceTimestampMode::None => return Ok(()),
148    };
149
150    socket.so_timestamping(options, bind_phc.unwrap_or_default())
151}
152
153pub fn open_interface_udp(
154    interface: InterfaceName,
155    port: u16,
156    timestamping: InterfaceTimestampMode,
157    bind_phc: Option<u32>,
158) -> std::io::Result<Socket<SocketAddr, Open>> {
159    // Setup the socket
160    let socket = RawSocket::open(libc::PF_INET6, libc::SOCK_DGRAM, libc::IPPROTO_UDP)?;
161    socket.enable_destination_ipv4()?;
162    socket.enable_destination_ipv6()?;
163    socket.reuse_addr()?;
164    socket.ipv6_v6only(false)?;
165    socket.bind(SocketAddrV6::new(Ipv6Addr::UNSPECIFIED, port, 0, 0).to_sockaddr(PrivateToken))?;
166    socket.bind_to_device(interface)?;
167    socket.ipv6_multicast_if(interface)?;
168    socket.ipv6_multicast_loop(false)?;
169    configure_timestamping(&socket, Some(interface), timestamping, bind_phc)?;
170    match timestamping {
171        InterfaceTimestampMode::HardwareAll | InterfaceTimestampMode::HardwareRecv => {
172            socket.driver_enable_hardware_timestamping(interface, libc::HWTSTAMP_FILTER_ALL as _)?
173        }
174        InterfaceTimestampMode::HardwarePTPAll | InterfaceTimestampMode::HardwarePTPRecv => socket
175            .driver_enable_hardware_timestamping(
176                interface,
177                libc::HWTSTAMP_FILTER_PTP_V2_L4_EVENT as _,
178            )?,
179        InterfaceTimestampMode::None
180        | InterfaceTimestampMode::SoftwareAll
181        | InterfaceTimestampMode::SoftwareRecv => {}
182    }
183    socket.set_nonblocking(true)?;
184
185    let local_addr = SocketAddr::from_sockaddr(socket.getsockname()?, PrivateToken)
186        .ok_or::<std::io::Error>(std::io::ErrorKind::Other.into())?;
187
188    Ok(Socket {
189        timestamp_mode: timestamping,
190        socket: Arc::new(AsyncFd::new(socket)?),
191        send_counter: std::sync::Mutex::new(0),
192        local_addr,
193        _state: PhantomData,
194    })
195}
196
197pub fn open_interface_udp4(
198    interface: InterfaceName,
199    port: u16,
200    timestamping: InterfaceTimestampMode,
201    bind_phc: Option<u32>,
202) -> std::io::Result<Socket<SocketAddrV4, Open>> {
203    // Setup the socket
204    let socket = RawSocket::open(libc::PF_INET, libc::SOCK_DGRAM, libc::IPPROTO_UDP)?;
205    socket.enable_destination_ipv4()?;
206    socket.reuse_addr()?;
207    socket.bind(SocketAddrV4::new(Ipv4Addr::UNSPECIFIED, port).to_sockaddr(PrivateToken))?;
208    socket.bind_to_device(interface)?;
209    socket.ip_multicast_if(interface)?;
210    socket.ip_multicast_loop(false)?;
211    configure_timestamping(&socket, Some(interface), timestamping, bind_phc)?;
212    match timestamping {
213        InterfaceTimestampMode::HardwareAll | InterfaceTimestampMode::HardwareRecv => {
214            socket.driver_enable_hardware_timestamping(interface, libc::HWTSTAMP_FILTER_ALL as _)?
215        }
216        InterfaceTimestampMode::HardwarePTPAll | InterfaceTimestampMode::HardwarePTPRecv => socket
217            .driver_enable_hardware_timestamping(
218                interface,
219                libc::HWTSTAMP_FILTER_PTP_V2_L4_EVENT as _,
220            )?,
221        InterfaceTimestampMode::None
222        | InterfaceTimestampMode::SoftwareAll
223        | InterfaceTimestampMode::SoftwareRecv => {}
224    }
225    socket.set_nonblocking(true)?;
226
227    let local_addr = SocketAddrV4::from_sockaddr(socket.getsockname()?, PrivateToken)
228        .ok_or::<std::io::Error>(std::io::ErrorKind::Other.into())?;
229
230    Ok(Socket {
231        timestamp_mode: timestamping,
232        socket: Arc::new(AsyncFd::new(socket)?),
233        send_counter: std::sync::Mutex::new(0),
234        local_addr,
235        _state: PhantomData,
236    })
237}
238
239pub fn open_interface_udp6(
240    interface: InterfaceName,
241    port: u16,
242    timestamping: InterfaceTimestampMode,
243    bind_phc: Option<u32>,
244) -> std::io::Result<Socket<SocketAddrV6, Open>> {
245    // Setup the socket
246    let socket = RawSocket::open(libc::PF_INET6, libc::SOCK_DGRAM, libc::IPPROTO_UDP)?;
247    socket.enable_destination_ipv6()?;
248    socket.reuse_addr()?;
249    socket.ipv6_v6only(true)?;
250    socket.bind(SocketAddrV6::new(Ipv6Addr::UNSPECIFIED, port, 0, 0).to_sockaddr(PrivateToken))?;
251    socket.bind_to_device(interface)?;
252    socket.ipv6_multicast_if(interface)?;
253    socket.ipv6_multicast_loop(false)?;
254    configure_timestamping(&socket, Some(interface), timestamping, bind_phc)?;
255    match timestamping {
256        InterfaceTimestampMode::HardwareAll | InterfaceTimestampMode::HardwareRecv => {
257            socket.driver_enable_hardware_timestamping(interface, libc::HWTSTAMP_FILTER_ALL as _)?
258        }
259        InterfaceTimestampMode::HardwarePTPAll | InterfaceTimestampMode::HardwarePTPRecv => socket
260            .driver_enable_hardware_timestamping(
261                interface,
262                libc::HWTSTAMP_FILTER_PTP_V2_L4_EVENT as _,
263            )?,
264        InterfaceTimestampMode::None
265        | InterfaceTimestampMode::SoftwareAll
266        | InterfaceTimestampMode::SoftwareRecv => {}
267    }
268    socket.set_nonblocking(true)?;
269
270    let local_addr = SocketAddrV6::from_sockaddr(socket.getsockname()?, PrivateToken)
271        .ok_or::<std::io::Error>(std::io::ErrorKind::Other.into())?;
272
273    Ok(Socket {
274        timestamp_mode: timestamping,
275        socket: Arc::new(AsyncFd::new(socket)?),
276        send_counter: std::sync::Mutex::new(0),
277        local_addr,
278        _state: PhantomData,
279    })
280}
281
282pub fn open_interface_ethernet(
283    interface: InterfaceName,
284    protocol: u16,
285    timestamping: InterfaceTimestampMode,
286    bind_phc: Option<u32>,
287) -> std::io::Result<Socket<EthernetAddress, Open>> {
288    let socket = RawSocket::open(
289        libc::AF_PACKET,
290        libc::SOCK_DGRAM,
291        u16::from_ne_bytes(protocol.to_be_bytes()) as _,
292    )?;
293    socket.bind(
294        EthernetAddress::new(
295            u16::from_ne_bytes(protocol.to_le_bytes()),
296            MacAddress::new([0; 6]),
297            interface
298                .get_index()
299                .ok_or(std::io::ErrorKind::InvalidInput)? as _,
300        )
301        .to_sockaddr(PrivateToken),
302    )?;
303    configure_timestamping(&socket, Some(interface), timestamping, bind_phc)?;
304    match timestamping {
305        InterfaceTimestampMode::HardwareAll | InterfaceTimestampMode::HardwareRecv => {
306            socket.driver_enable_hardware_timestamping(interface, libc::HWTSTAMP_FILTER_ALL as _)?
307        }
308        InterfaceTimestampMode::HardwarePTPAll | InterfaceTimestampMode::HardwarePTPRecv => socket
309            .driver_enable_hardware_timestamping(
310                interface,
311                libc::HWTSTAMP_FILTER_PTP_V2_L2_EVENT as _,
312            )?,
313        InterfaceTimestampMode::None
314        | InterfaceTimestampMode::SoftwareAll
315        | InterfaceTimestampMode::SoftwareRecv => {}
316    }
317    socket.set_nonblocking(true)?;
318
319    let local_addr = EthernetAddress::from_sockaddr(socket.getsockname()?, PrivateToken)
320        .ok_or::<std::io::Error>(std::io::ErrorKind::Other.into())?;
321
322    Ok(Socket {
323        timestamp_mode: timestamping,
324        socket: Arc::new(AsyncFd::new(socket)?),
325        send_counter: std::sync::Mutex::new(0),
326        local_addr,
327        _state: PhantomData,
328    })
329}
330
331#[cfg(test)]
332mod tests {
333    use std::net::IpAddr;
334
335    use crate::socket::{connect_address, open_ip, GeneralTimestampMode};
336
337    use super::*;
338
339    #[tokio::test]
340    async fn test_open_udp() {
341        use std::str::FromStr;
342        let a = open_interface_udp(
343            InterfaceName::from_str("lo").unwrap(),
344            5128,
345            super::InterfaceTimestampMode::None,
346            None,
347        )
348        .unwrap();
349
350        let mut b = connect_address(
351            SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 5128),
352            GeneralTimestampMode::None,
353        )
354        .unwrap();
355        assert!(b.send(&[1, 2, 3]).await.is_ok());
356        let mut buf = [0; 4];
357        let recv_result = a.recv(&mut buf).await.unwrap();
358        assert_eq!(recv_result.bytes_read, 3);
359        assert_eq!(&buf[0..3], &[1, 2, 3]);
360        assert_eq!(
361            recv_result.local_addr,
362            SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 5128)
363        );
364
365        let mut b = connect_address(
366            SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 1)), 5128),
367            GeneralTimestampMode::None,
368        )
369        .unwrap();
370        assert!(b.send(&[1, 2, 3]).await.is_ok());
371        let mut buf = [0; 4];
372        let recv_result = a.recv(&mut buf).await.unwrap();
373        assert_eq!(recv_result.bytes_read, 3);
374        assert_eq!(&buf[0..3], &[1, 2, 3]);
375        assert_eq!(
376            recv_result.local_addr,
377            SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 1, 1)), 5128)
378        );
379    }
380
381    #[tokio::test]
382    async fn test_open_ip_reuse_addr_after_interface() {
383        use std::str::FromStr;
384        let _a = open_interface_udp(
385            InterfaceName::from_str("lo").unwrap(),
386            5132,
387            super::InterfaceTimestampMode::None,
388            None,
389        )
390        .unwrap();
391        let _b = open_ip(
392            SocketAddr::new(Ipv4Addr::UNSPECIFIED.into(), 5132),
393            GeneralTimestampMode::None,
394            true,
395        )
396        .unwrap();
397    }
398
399    #[tokio::test]
400    async fn test_open_ip_reuse_addr_before_interface() {
401        use std::str::FromStr;
402        let _a = open_ip(
403            SocketAddr::new(Ipv4Addr::UNSPECIFIED.into(), 5133),
404            GeneralTimestampMode::None,
405            true,
406        )
407        .unwrap();
408        let _b = open_interface_udp(
409            InterfaceName::from_str("lo").unwrap(),
410            5133,
411            super::InterfaceTimestampMode::None,
412            None,
413        )
414        .unwrap();
415    }
416
417    #[tokio::test]
418    async fn test_open_udp6() {
419        use std::str::FromStr;
420        let mut a = open_interface_udp6(
421            InterfaceName::from_str("lo").unwrap(),
422            5123,
423            super::InterfaceTimestampMode::None,
424            None,
425        )
426        .unwrap();
427        let mut b = connect_address(
428            SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 5123),
429            GeneralTimestampMode::None,
430        )
431        .unwrap();
432        assert!(b.send(&[1, 2, 3]).await.is_ok());
433        let mut buf = [0; 4];
434        let recv_result = a.recv(&mut buf).await.unwrap();
435        assert_eq!(recv_result.bytes_read, 3);
436        assert_eq!(&buf[0..3], &[1, 2, 3]);
437        assert_eq!(
438            recv_result.local_addr,
439            SocketAddrV6::new(Ipv6Addr::LOCALHOST, 5123, 0, 0)
440        );
441        assert!(a.send_to(&[4, 5, 6], recv_result.remote_addr).await.is_ok());
442        let recv_result = b.recv(&mut buf).await.unwrap();
443        assert_eq!(recv_result.bytes_read, 3);
444        assert_eq!(&buf[0..3], &[4, 5, 6]);
445    }
446
447    #[tokio::test]
448    async fn test_open_udp4() {
449        use std::str::FromStr;
450        let mut a = open_interface_udp4(
451            InterfaceName::from_str("lo").unwrap(),
452            5124,
453            super::InterfaceTimestampMode::None,
454            None,
455        )
456        .unwrap();
457        let mut b = connect_address(
458            SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 5124),
459            GeneralTimestampMode::None,
460        )
461        .unwrap();
462        assert!(b.send(&[1, 2, 3]).await.is_ok());
463        let mut buf = [0; 4];
464        let recv_result = a.recv(&mut buf).await.unwrap();
465        assert_eq!(recv_result.bytes_read, 3);
466        assert_eq!(&buf[0..3], &[1, 2, 3]);
467        assert_eq!(
468            recv_result.local_addr,
469            SocketAddrV4::new(Ipv4Addr::LOCALHOST, 5124)
470        );
471        assert!(a.send_to(&[4, 5, 6], recv_result.remote_addr).await.is_ok());
472        let recv_result = b.recv(&mut buf).await.unwrap();
473        assert_eq!(recv_result.bytes_read, 3);
474        assert_eq!(&buf[0..3], &[4, 5, 6]);
475    }
476
477    #[tokio::test]
478    async fn test_software_timestamping() {
479        use std::time::SystemTime;
480
481        let a = open_ip(
482            SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 5126),
483            GeneralTimestampMode::SoftwareAll,
484            false,
485        )
486        .unwrap();
487        let mut b = connect_address(
488            SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 5126),
489            GeneralTimestampMode::SoftwareAll,
490        )
491        .unwrap();
492
493        let before = SystemTime::now();
494        let send_ts = b
495            .send(&[1, 2, 3])
496            .await
497            .unwrap()
498            .selected_timestamp()
499            .unwrap();
500        let after = SystemTime::now();
501
502        let mut buf = [0; 4];
503        let recv_result = a.recv(&mut buf).await.unwrap();
504        let recv_ts = recv_result.timestamp_data.selected_timestamp().unwrap();
505
506        let before = before
507            .duration_since(SystemTime::UNIX_EPOCH)
508            .unwrap()
509            .as_secs();
510        let after = after
511            .duration_since(SystemTime::UNIX_EPOCH)
512            .unwrap()
513            .as_secs();
514        assert!((send_ts.seconds - (before as i64)).abs() < 2);
515        assert!((send_ts.seconds - (after as i64)).abs() < 2);
516
517        let send_nanos = send_ts.seconds * 1_000_000_000 + (send_ts.nanos as i64);
518        let recv_nanos = recv_ts.seconds * 1_000_000_000 + (recv_ts.nanos as i64);
519        assert!((send_nanos - recv_nanos) < 1_000_000 * 10);
520    }
521}