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 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
37fn 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 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 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 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 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 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 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}