//! The UDP receive loop. //! //! `SO_REUSEPORT` lets several tasks bind the *same* port and have the kernel //! spread datagrams across them, which is how one port scales past one core //! without a dispatcher task in the middle. Each worker owns its own socket, so //! there is no shared state on the hot path at all. //! //! The loop deliberately contains no `.await` between receiving and replying: the //! whole of [`Ingest::handle`] is synchronous, and a full writer channel is //! answered rather than waited on. Applying backpressure here would turn one slow //! disk into packet loss for every device at once. use std::net::SocketAddr; use std::sync::Arc; use anyhow::{Context, Result}; use socket2::{Domain, Protocol, Socket, Type}; use tokio::net::UdpSocket; use tracing::{error, info, warn}; use crate::db::now; use crate::ingest::{Ingest, Peer, Transport}; /// Receive buffer per socket. A queue flush from every device at once arrives as /// a burst, and the kernel's default (a few hundred kB) drops it on the floor /// before this process ever sees it. const RECV_BUFFER_BYTES: usize = 2 * 1024 * 1024; /// One more than the largest datagram, so an oversized packet is *seen* to be /// oversized instead of being silently truncated into something that might parse. const READ_BUFFER: usize = otproto::MAX_DATAGRAM + 1; fn bind_reuseport(addr: SocketAddr) -> Result { let domain = if addr.is_ipv6() { Domain::IPV6 } else { Domain::IPV4 }; let socket = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP)).context("creating socket")?; socket.set_reuse_address(true).context("SO_REUSEADDR")?; socket.set_reuse_port(true).context("SO_REUSEPORT")?; // Best effort: on Linux the kernel doubles the requested value and caps it at // net.core.rmem_max, so a smaller buffer than asked for is normal and not // worth failing startup over. if let Err(e) = socket.set_recv_buffer_size(RECV_BUFFER_BYTES) { warn!(error = %e, "could not enlarge the UDP receive buffer; bursts may be dropped"); } socket.set_nonblocking(true).context("set_nonblocking")?; socket .bind(&addr.into()) .with_context(|| format!("binding {addr}/udp"))?; UdpSocket::from_std(socket.into()).context("handing the socket to tokio") } /// Spawn `workers` receive tasks on `addr`. pub fn spawn( addr: SocketAddr, workers: usize, ingest: Arc, ) -> Result>> { let mut tasks = Vec::with_capacity(workers); for id in 0..workers { let socket = bind_reuseport(addr)?; tasks.push(tokio::spawn(run(id, socket, Arc::clone(&ingest)))); } info!(%addr, workers, "OTP/1 UDP listener started"); // The trap worth stating in the log, because operators hit it on day one and // the symptom (everything falls back to TLS) is far from the cause. info!("reminder: HTTP reverse proxies do not forward UDP — {addr} needs its own firewall rule"); Ok(tasks) } async fn run(id: usize, socket: UdpSocket, ingest: Arc) { let mut buf = vec![0u8; READ_BUFFER]; loop { let (len, peer_addr) = match socket.recv_from(&mut buf).await { Ok(v) => v, Err(e) => { // On UDP a send error can surface here as ICMP-driven // ECONNREFUSED for a *previous* send. It says nothing about the // socket's health, so log and keep going rather than exiting the // worker and silently losing a quarter of the capacity. warn!(worker = id, error = %e, "recv_from failed"); continue; } }; let peer = Peer { addr: peer_addr, transport: Transport::Udp, }; if let Some(reply) = ingest.handle(&buf[..len], peer, now()) && let Err(e) = socket.send_to(&reply, peer_addr).await { // The phone will retry; there is nothing to recover here. error!(worker = id, %peer_addr, error = %e, "sending reply failed"); } } } #[cfg(test)] mod tests { use super::*; use otproto::msg::Direction; use otproto::{Key, Message, Point, kdf}; use std::time::Duration; const TOKEN_ID: u64 = 0xABCD_0123_4567_89EF; const TOKEN_KEY: Key = [0x33; 32]; /// Round trip a real datagram over a real loopback socket. This is the only /// test that exercises the socket options and the reply path together. #[tokio::test] async fn a_datagram_over_loopback_is_acked() { let dir = tempfile::tempdir().expect("temp dir"); let db = crate::db::Db::open(&dir.path().join("t.db")) .await .expect("open"); sqlx::query( "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \ VALUES (1, 'a', 'x', 'A', 0, 0)", ) .execute(&db.write) .await .expect("user"); let (writer, _task) = crate::writer::spawn(db.write.clone()); let ingest = Arc::new(Ingest::new(writer, 30 * 86_400, None)); ingest.insert_token(crate::ingest::TokenSlot::new(TOKEN_ID, 1, &TOKEN_KEY, 1)); // Port 0 lets the OS choose; then read it back, because SO_REUSEPORT // workers must all bind the *same* concrete port. let probe = bind_reuseport("127.0.0.1:0".parse().expect("literal")).expect("bind"); let addr = probe.local_addr().expect("local addr"); drop(probe); let _tasks = spawn(addr, 2, Arc::clone(&ingest)).expect("spawn"); let client = UdpSocket::bind("127.0.0.1:0").await.expect("client bind"); let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); let msg = Message::Loc(vec![Point { acc_dm: Some(50), ..Point::new(now() as u32, 525_200_080, 134_050_000) }]); let datagram = otproto::seal_message(&k_up, TOKEN_ID, [0x77; 12], &msg); client.send_to(&datagram, addr).await.expect("send"); let mut buf = vec![0u8; READ_BUFFER]; let len = tokio::time::timeout(Duration::from_secs(2), client.recv(&mut buf)) .await .expect("no reply within 2s") .expect("recv"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); match otproto::open_message(&k_down, &buf[..len]) .expect("ack opens") .1 { Message::Ack(ack) => assert_eq!(ack.nonces, vec![[0x77u8; 12]]), other => panic!("expected an ACK, got {other:?}"), } assert!(len <= datagram.len(), "the reply amplified the request"); } #[tokio::test] async fn several_workers_can_share_one_port() { let probe = bind_reuseport("127.0.0.1:0".parse().expect("literal")).expect("first bind"); let addr = probe.local_addr().expect("local addr"); // The second bind on the same concrete port is the thing SO_REUSEPORT // makes legal, and the thing the whole multi-worker design rests on. let second = bind_reuseport(addr).expect("SO_REUSEPORT should allow a second bind"); assert_eq!(second.local_addr().expect("addr"), addr); } }