//! Transport-agnostic datagram handling. //! //! [`Ingest::handle`] is **synchronous** and takes a couple of microseconds. It //! never touches the database: every active token lives in an in-memory map, //! loaded at startup and updated on mint/revoke. That is what lets the UDP //! receive loop stay a tight `recv_from` → `handle` → `send_to` cycle with no //! `.await` in the middle, and it is why an unknown `token_id` is genuinely //! unknown rather than merely uncached. //! //! The cost is memory proportional to the number of active tokens. At ~90 bytes //! per slot, a million tokens would be 90 MB; for a self-hosted instance with a //! handful of users it is a few kilobytes. If that ever stops being true, this is //! the module to revisit. //! //! The same function serves the UDP loop and the TLS fallback listener, so the //! two transports cannot drift apart in their handling of anything. use std::net::{IpAddr, SocketAddr}; use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use dashmap::DashMap; #[cfg(test)] use otproto::msg::Direction; use otproto::{ Ack, AckFlags, DecodeError, Header, Key, MAX_POINTS, Message, Nack, NackReason, Point, RevokeReason, Revoked, kdf, }; use rand::TryRngCore; use rand::rngs::OsRng; use tracing::{debug, trace}; use crate::limits::Limits; use crate::writer::{Accepted, WriteHandle, WriteOp}; /// Which transport a datagram arrived on. Recorded per token so the UI can show /// whether a phone is on UDP or has fallen back to TLS. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Transport { Udp, /// Constructed by the TLS-over-TCP fallback listener, which lands in a later /// step; the recording path is already transport-agnostic so that listener /// only has to call `Ingest::handle` with this variant. #[allow(dead_code)] Tls, } impl Transport { const fn as_str(self) -> &'static str { match self { Self::Udp => "udp", Self::Tls => "tls", } } } #[derive(Debug, Clone, Copy)] pub struct Peer { pub addr: SocketAddr, pub transport: Transport, } impl Peer { pub fn ip(&self) -> IpAddr { self.addr.ip() } } /// Everything needed to serve one token, with no database round trip. #[derive(Debug)] pub struct TokenSlot { pub token_id: u64, /// Resolved once, at load time, so the write path never has to ask which /// account a token belongs to. pub user_id: i64, pub k_up: Key, pub k_down: Key, pub config_version: u16, /// Revoked tokens stay in the map instead of being removed. /// /// The key is what lets the server answer at all, and a revoked device is /// precisely the one it needs to answer. Dropping the slot would leave the /// only possible reply an unauthenticated one — a reflector — where keeping /// it makes the reply a sealed `NACK` nobody else could have produced. /// /// The token authorises nothing from the moment this is set; the key survives /// only so the server can say so. pub revoked: Option, } impl TokenSlot { pub fn new(token_id: u64, user_id: i64, token_key: &Key, config_version: u16) -> Self { let (k_up, k_down) = kdf::derive_both(token_key); Self { token_id, user_id, k_up, k_down, config_version, revoked: None, } } #[must_use] pub fn revoked(mut self, reason: RevokeReason) -> Self { self.revoked = Some(reason); self } } /// Aggregate counters. Per-packet logging would itself be an amplifier, so these /// are the only per-packet observability, exposed on a localhost-only `/metrics`. #[derive(Debug, Default)] pub struct Counters { pub received: AtomicU64, pub malformed: AtomicU64, pub unknown_token: AtomicU64, pub auth_failed: AtomicU64, pub rate_limited: AtomicU64, pub throttled: AtomicU64, pub points_accepted: AtomicU64, pub points_rejected: AtomicU64, pub acks_sent: AtomicU64, pub nacks_sent: AtomicU64, pub silent_drops: AtomicU64, /// Sealed notices to a token the server still holds a key for. pub revoked_notices_sent: AtomicU64, /// Notices sent without being able to verify the request. The reflection /// budget, in other words — worth watching. pub unverified_notices_sent: AtomicU64, pub notices_suppressed: AtomicU64, } impl Counters { fn bump(counter: &AtomicU64) { counter.fetch_add(1, Ordering::Relaxed); } fn add(counter: &AtomicU64, n: u64) { counter.fetch_add(n, Ordering::Relaxed); } } pub struct Ingest { tokens: DashMap>, limits: Limits, writer: WriteHandle, pub counters: Counters, /// Accept timestamps within ±this many seconds of the server clock. ts_window_s: u32, /// Master for per-token revocation keys, or `None` when unverified notices /// are switched off in config. /// /// `Some` is the only thing that lets this server reply to a datagram it /// cannot verify, so the config switch is represented as the presence of the /// key rather than as a separate boolean. There is then no way to enable the /// behaviour by accident. revocation_master: Option, } impl Ingest { pub fn new(writer: WriteHandle, ts_window_s: u32, revocation_master: Option) -> Self { Self { tokens: DashMap::new(), limits: Limits::new(), writer, counters: Counters::default(), ts_window_s, revocation_master, } } pub fn insert_token(&self, slot: TokenSlot) { self.tokens.insert(slot.token_id, Arc::new(slot)); } /// Called by revoke, by password change, and by the staleness sweep. /// /// The slot stays, flagged. The token stops authorising anything /// immediately; what survives is only the ability to tell that one device it /// is finished, in a message nobody else could have sealed. pub fn mark_revoked(&self, token_id: u64, reason: RevokeReason) { if let Some(existing) = self.tokens.get(&token_id).map(|s| Arc::clone(&s)) { if existing.revoked.is_some() { return; } self.tokens.insert( token_id, Arc::new(TokenSlot { revoked: Some(reason), ..*existing }), ); } } /// Active tokens only — a revoked slot is bookkeeping, not a live device. pub fn active_token_count(&self) -> usize { self.tokens.iter().filter(|s| s.revoked.is_none()).count() } pub fn limits(&self) -> &Limits { &self.limits } /// Handle one datagram. Returns the bytes to send back, if any. /// /// `now` is passed in rather than read from the clock so this is testable /// without sleeping. pub fn handle(&self, datagram: &[u8], peer: Peer, now: i64) -> Option> { Counters::bump(&self.counters.received); // 1. Structure. No state, no allocation, no crypto. let header = match Header::peek(datagram) { Ok(h) => h, Err(e) => { Counters::bump(&self.counters.malformed); trace!(?e, "dropping malformed datagram"); return self.silent(); } }; // A downlink type arriving on the uplink is either a bug or someone // replaying our own traffic back at us. It can never be legitimate. if !header.msg_type.is_uplink() { Counters::bump(&self.counters.malformed); return self.silent(); } // 2. Per-IP budget and bans. if let Err(reason) = self.limits.check_ip(peer.ip()) { Counters::bump(&self.counters.rate_limited); trace!(?reason, "dropping rate-limited datagram"); return self.silent(); } // 3. Does this token exist? let Some(slot) = self.tokens.get(&header.token_id).map(|s| Arc::clone(&s)) else { Counters::bump(&self.counters.unknown_token); self.limits.note_unknown_token(peer.ip()); // No key, so no way to verify this datagram — which makes any reply a // reply to an address the sender merely claimed. See // [`Self::unverified_notice`] for what makes that tolerable. return self.unverified_notice(header.token_id, peer, datagram.len()); }; // 4. AEAD. The first expensive step, ~1 µs for a 62-byte packet. // // Everything below this line may answer; nothing above it ever does. // That is not a style rule, it is the anti-reflection defence. UDP source // addresses are trivially forged, and `token_id` travels in cleartext, so // anyone who has seen one datagram can name a valid token. If the server // replied before verifying, an attacker could spoof a victim's address // and have us send them a packet per junk datagram — laundering the // attacker's origin and firing at whatever rate they choose. // // Which is why the per-token budget is checked *after* this and not // before, even though that costs an AEAD open on every flooded packet. // The per-IP budget above absorbs the bulk at no crypto cost. let payload = match otproto::open(&slot.k_up, datagram) { Ok((_, payload)) => payload, Err(DecodeError::AuthFailed) => { Counters::bump(&self.counters.auth_failed); self.limits.note_aead_failure(peer.ip()); // Never answer this. The server cannot know who sent it, so a // reply would be both a forgery oracle and a reflector. return self.silent(); } Err(_) => { Counters::bump(&self.counters.malformed); return self.silent(); } }; // 5. Is this token still alive? Checked after AEAD, so the answer is a // sealed message the real device can trust and nobody else can forge. // This is the whole reason revoked slots keep their keys. if let Some(reason) = slot.revoked { Counters::bump(&self.counters.revoked_notices_sent); debug!(token_id = slot.token_id, ?reason, "revoked token reported"); return self.seal_reply( &slot, Message::Nack(Nack { nonce: header.nonce, reason: NackReason::UnknownToken, retry_after_s: 0, }), ); } // 6. Per-token budget. Authenticated, so this NACK reaches the device // that actually sent the datagram and nobody else. if self.limits.check_token(header.token_id).is_err() { Counters::bump(&self.counters.rate_limited); return self.nack(&slot, header.nonce, NackReason::RateLimited, 5); } let msg = match Message::decode_payload(header.msg_type, &payload) { Ok(m) => m, Err(e) => { Counters::bump(&self.counters.malformed); debug!( token_id = header.token_id, ?e, "authenticated but malformed payload" ); return self.nack(&slot, header.nonce, NackReason::Malformed, 0); } }; // From here the packet is authenticated, so telemetry is safe to record. self.record_seen(&slot, peer, now); match msg { Message::Loc(points) => self.handle_loc(&slot, header, points, now), Message::Hello(hello) => { self.writer.try_send(WriteOp::TokenHello { token_id: slot.token_id as i64, app_version: i64::from(hello.app_version_code), os_api_level: i64::from(hello.os_api_level), }); self.ack(&slot, vec![header.nonce], hello.config_version) } Message::ConfigGet(_) => { // CONFIG is served by the HTTP/state path in this build; a device // asking gets silence rather than a stale answer, and retries. // Wired up with the config editor in a later step. self.silent() } Message::Ping(ping) => { // The echo is opaque to us; copying it back is the whole job. let pong = Message::Pong(otproto::Pong { echo: ping.echo, seq: ping.seq, }); self.seal_reply(&slot, pong) } // Unreachable: the direction check above rejected every downlink type. Message::Ack(_) | Message::Nack(_) | Message::Config(_) | Message::Pong(_) | Message::Revoked(_) => { Counters::bump(&self.counters.malformed); self.silent() } } } fn handle_loc( &self, slot: &TokenSlot, header: Header, points: Vec, now: i64, ) -> Option> { // Semantic validation. A point far outside the timestamp window would // land where the retention sweep never reaches it, so it is dropped // rather than stored — but the rest of the batch is still kept, because // one bad fix must not cost the user a whole journey. let before = points.len(); let accepted: Vec = points .into_iter() .filter(|p| p.validate(now as u32, self.ts_window_s).is_ok()) .collect(); let rejected = before - accepted.len(); Counters::add(&self.counters.points_rejected, rejected as u64); if accepted.is_empty() { Counters::bump(&self.counters.malformed); return self.nack(slot, header.nonce, NackReason::Malformed, 0); } Counters::add(&self.counters.points_accepted, accepted.len() as u64); let accepted_count = accepted.len(); match self.writer.try_send(WriteOp::Points { user_id: slot.user_id, src_token_id: slot.token_id as i64, points: accepted, recv_at: now, }) { Accepted::Yes => {} Accepted::Saturated => { // Do not ack what we did not store: the client keeps the points // queued and retries. THROTTLE tells it to slow down first. Counters::bump(&self.counters.throttled); return self.throttled(slot, header.nonce); } Accepted::Closed => return self.silent(), } trace!( token_id = slot.token_id, user_id = slot.user_id, accepted_count, "stored points" ); self.ack(slot, vec![header.nonce], slot.config_version) } fn record_seen(&self, slot: &TokenSlot, peer: Peer, now: i64) { // Best effort: if the writer is saturated, telemetry is the first thing // worth dropping. self.writer.try_send(WriteOp::TokenSeen { token_id: slot.token_id as i64, at: now, src_ip: peer.ip().to_string(), src_port: peer.addr.port(), transport: peer.transport.as_str(), }); } fn ack( &self, slot: &TokenSlot, nonces: Vec<[u8; 12]>, device_config_version: u16, ) -> Option> { let flags = if device_config_version < slot.config_version { AckFlags::CONFIG_PENDING } else { AckFlags::NONE }; debug_assert!(nonces.len() <= MAX_POINTS); let ack = Message::Ack(Ack { nonces, flags }); Counters::bump(&self.counters.acks_sent); self.seal_reply(slot, ack) } /// The writer channel is full, so the points were not stored. /// /// This is a `NACK`, not an `ACK` with the `THROTTLE` flag, because an `ACK` /// retires the nonces it names: acking here would tell the client to delete /// points that never reached the database. `RateLimited` with a retry hint is /// exactly the "keep them and back off" semantic the client already /// implements, so saturation costs a delay rather than data. fn throttled(&self, slot: &TokenSlot, nonce: [u8; 12]) -> Option> { self.nack(slot, nonce, NackReason::RateLimited, 10) } fn nack( &self, slot: &TokenSlot, nonce: [u8; 12], reason: NackReason, retry_after_s: u8, ) -> Option> { Counters::bump(&self.counters.nacks_sent); self.seal_reply( slot, Message::Nack(Nack { nonce, reason, retry_after_s, }), ) } /// The one reply this server sends without having verified the request. /// /// The server has no record of this `token_id`, so it cannot open the /// datagram and cannot know who really sent it. `peer` is whatever the source /// address claimed, which on UDP is forgeable. Answering therefore makes this /// port a reflector, and three things are what keep that from mattering: /// /// 1. **It cannot amplify.** [`Revoked`] is 38 bytes, the smallest message in /// the protocol, and a request shorter than that gets nothing. /// 2. **It is rate limited by destination.** The budget is keyed on the /// address the reply would go to — the victim, for a spoofed packet — at /// one per minute, under a global ceiling for distributed attempts. /// 3. **It is optional.** No revocation master configured, no reply. /// /// It cannot be forged, which is the other half. `K_rev` is derived from a /// server master and this `token_id`, so a third party cannot produce one and /// neither can another legitimate device — each only ever learns its own. /// /// This path exists for the cases where the row is genuinely gone: a database /// restored from a backup that predates the login, or a rotated server key. /// The ordinary revocation path keeps the slot and answers with a sealed /// `NACK`, which needs none of the above. fn unverified_notice(&self, token_id: u64, peer: Peer, request_len: usize) -> Option> { let Some(master) = self.revocation_master else { return self.silent(); }; if !otproto::may_answer_unverified(request_len) { return self.silent(); } if self.limits.check_notice(peer.ip()).is_err() { Counters::bump(&self.counters.notices_suppressed); return self.silent(); } let mut nonce = [0u8; 12]; if OsRng.try_fill_bytes(&mut nonce).is_err() { return self.silent(); } let msg = Message::Revoked(Revoked { reason: RevokeReason::Unknown, }); let reply = otproto::seal_message( &otproto::revocation_key(&master, token_id), token_id, nonce, &msg, ); debug_assert!( reply.len() <= request_len, "an unverified reply must never exceed its request" ); Counters::bump(&self.counters.unverified_notices_sent); Some(reply) } /// Seal a reply. /// /// Anti-amplification is not enforced here, because by this point it is /// already established: a reply is only reached after the datagram passed AEAD /// verification, so a reflection attacker must hold a live token key, and /// every reply this server can construct fits in [`otproto::MAX_REPLY`] bytes /// — a ratio near 1 against any request. The `debug_assert` is a development /// guard against a future message type outgrowing that budget, not a runtime /// control. fn seal_reply(&self, slot: &TokenSlot, msg: Message) -> Option> { debug_assert!( otproto::fits_reply_budget(otproto::datagram_len(msg.payload_len())), "{:?} exceeds the reply budget of {} bytes", msg.msg_type(), otproto::MAX_REPLY, ); let mut nonce = [0u8; 12]; if OsRng.try_fill_bytes(&mut nonce).is_err() { // Without fresh randomness we must not seal anything: reusing a nonce // under ChaCha20-Poly1305 leaks the keystream. return self.silent(); } Some(otproto::seal_message( &slot.k_down, slot.token_id, nonce, &msg, )) } fn silent(&self) -> Option> { Counters::bump(&self.counters.silent_drops); None } /// Derived key for a direction. Test-only: the live paths look up the slot and /// use both keys, and exposing a key by id elsewhere would invite misuse. #[cfg(test)] pub fn key_for(&self, token_id: u64, dir: Direction) -> Option { self.tokens.get(&token_id).map(|s| match dir { Direction::Up => s.k_up, Direction::Down => s.k_down, }) } } #[cfg(test)] mod tests { use super::*; use crate::db::Db; use otproto::point::Flags; use otproto::{HEADER_LEN, TAG_LEN}; const TOKEN_ID: u64 = 0x0123_4567_89AB_CDEF; const TOKEN_KEY: Key = [0x5A; 32]; const REVOCATION_MASTER: Key = [0xC3; 32]; const NOW: i64 = 1_785_000_042; fn peer() -> Peer { Peer { addr: "203.0.113.5:40000".parse().expect("literal"), transport: Transport::Udp, } } async fn fixture() -> (Ingest, Db, tempfile::TempDir, tokio::task::JoinHandle<()>) { let dir = tempfile::tempdir().expect("temp dir"); let 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 = Ingest::new(writer, 30 * 86_400, Some(REVOCATION_MASTER)); ingest.insert_token(TokenSlot::new(TOKEN_ID, 1, &TOKEN_KEY, 1)); (ingest, db, dir, task) } fn seal(msg: &Message, nonce_seed: u8) -> Vec { let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); otproto::seal_message(&k_up, TOKEN_ID, [nonce_seed; 12], msg) } fn loc(ts: u32) -> Message { Message::Loc(vec![Point { acc_dm: Some(80), flags: Flags::NONE, ..Point::new(ts, 525_200_080, 134_050_000) }]) } /// The anti-reflection invariant, stated as a test: nothing that fails the /// AEAD check may produce a reply. /// /// UDP source addresses are forgeable, so any reply to an unverified /// datagram is a packet an attacker can aim at a third party. `token_id` is /// cleartext, so naming a real token costs nothing — which is what makes the /// per-token rate limit the interesting case here rather than a theoretical /// one. Its budget is small enough that a flood trips it immediately. #[tokio::test] async fn nothing_that_fails_aead_ever_gets_a_reply() { let (ingest, _db, _dir, _task) = fixture().await; let valid = seal(&loc(NOW as u32), 1); // Far past any per-token budget, so the pre-AEAD ordering bug would show // up here as a rate-limit NACK sent to an unauthenticated sender. for i in 0..200 { // A real header naming a real token, with a corrupted tag. let mut forged = valid.clone(); let last = forged.len() - 1; forged[last] ^= 1; forged[HEADER_LEN] ^= i as u8; assert_eq!( ingest.handle(&forged, peer(), NOW), None, "a datagram that fails AEAD was answered on attempt {i}" ); // And a header-only datagram, which cannot authenticate at all. let stub = valid[..HEADER_LEN + TAG_LEN].to_vec(); assert_eq!( ingest.handle(&stub, peer(), NOW), None, "a truncated datagram was answered on attempt {i}" ); } assert_eq!( ingest.counters.acks_sent.load(Ordering::Relaxed) + ingest.counters.nacks_sent.load(Ordering::Relaxed), 0, "the server sent something in response to unauthenticated traffic" ); } #[tokio::test] async fn a_valid_loc_is_acked_and_stored() { let (ingest, db, _dir, _task) = fixture().await; let datagram = seal(&loc(NOW as u32), 1); let reply = ingest.handle(&datagram, peer(), NOW).expect("should ack"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); let (_, msg) = otproto::open_message(&k_down, &reply).expect("ack opens"); match msg { Message::Ack(ack) => { assert_eq!( ack.nonces, vec![[1u8; 12]], "the ack must echo the request nonce" ); } other => panic!("expected an ACK, got {other:?}"), } tokio::time::sleep(std::time::Duration::from_millis(400)).await; let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points WHERE user_id = 1") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 1); } #[tokio::test] async fn an_unknown_token_is_counted_and_struck() { let (ingest, _db, _dir, _task) = fixture().await; let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); let datagram = otproto::seal_message(&k_up, 0xDEAD_BEEF, [2; 12], &loc(NOW as u32)); // The notice itself is covered by // `an_unknown_token_draws_one_small_rate_limited_notice`; what matters // here is that answering did not stop the sender being treated as a // scanner. Naming unknown ids still earns strikes and eventually a ban. let _ = ingest.handle(&datagram, peer(), NOW); assert_eq!(ingest.counters.unknown_token.load(Ordering::Relaxed), 1); assert_eq!(ingest.counters.points_accepted.load(Ordering::Relaxed), 0); } #[tokio::test] async fn a_forged_datagram_gets_silence() { let (ingest, _db, _dir, _task) = fixture().await; let mut datagram = seal(&loc(NOW as u32), 3); let last = datagram.len() - 1; datagram[last] ^= 1; assert!( ingest.handle(&datagram, peer(), NOW).is_none(), "a failed tag must never be answered" ); assert_eq!(ingest.counters.auth_failed.load(Ordering::Relaxed), 1); } #[tokio::test] async fn a_downlink_message_arriving_on_the_uplink_is_dropped() { let (ingest, _db, _dir, _task) = fixture().await; let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); let pong = Message::Pong(otproto::Pong { echo: 1, seq: 3 }); let datagram = otproto::seal_message(&k_down, TOKEN_ID, [4; 12], &pong); assert!(ingest.handle(&datagram, peer(), NOW).is_none()); } #[tokio::test] async fn a_ping_is_answered_with_a_pong_of_no_greater_size() { let (ingest, _db, _dir, _task) = fixture().await; let ping = Message::Ping(otproto::Ping { echo: 0xDEAD_BEEF, seq: 7, }); let datagram = seal(&ping, 5); let reply = ingest.handle(&datagram, peer(), NOW).expect("pong"); assert!(reply.len() <= datagram.len(), "PONG amplified the PING"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); match otproto::open_message(&k_down, &reply).expect("opens").1 { Message::Pong(p) => { assert_eq!(p.seq, 7); assert_eq!( p.echo, 0xDEAD_BEEF, "the PING's opaque echo must come back untouched" ); } other => panic!("expected a PONG, got {other:?}"), } } #[tokio::test] async fn a_point_outside_the_timestamp_window_is_rejected() { let (ingest, db, _dir, _task) = fixture().await; // Year 2100: far beyond anything the retention sweep would ever reach. let datagram = seal(&loc(4_102_444_800), 6); let reply = ingest.handle(&datagram, peer(), NOW).expect("nack"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); match otproto::open_message(&k_down, &reply).expect("opens").1 { Message::Nack(n) => assert_eq!(n.reason, NackReason::Malformed), other => panic!("expected a NACK, got {other:?}"), } tokio::time::sleep(std::time::Duration::from_millis(300)).await; let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 0); } #[tokio::test] async fn one_bad_point_does_not_cost_the_whole_batch() { let (ingest, db, _dir, _task) = fixture().await; let good = Point { acc_dm: Some(80), ..Point::new(NOW as u32, 525_200_080, 134_050_000) }; let bad = Point::new(NOW as u32, 910_000_000, 0); // impossible latitude let datagram = seal(&Message::Loc(vec![good, bad]), 7); assert!( ingest.handle(&datagram, peer(), NOW).is_some(), "should still ack" ); tokio::time::sleep(std::time::Duration::from_millis(400)).await; let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 1, "the good point must survive its bad neighbour"); assert_eq!(ingest.counters.points_rejected.load(Ordering::Relaxed), 1); } /// A revoked token keeps its key so this answer is possible at all. /// /// The alternative — dropping the slot — leaves silence as the only safe /// reply, and the phone keeps reporting into nothing until someone notices. #[tokio::test] async fn a_revoked_token_is_told_so_in_a_message_it_can_verify() { let (ingest, _db, _dir, _task) = fixture().await; ingest.mark_revoked(TOKEN_ID, RevokeReason::Revoked); let datagram = seal(&loc(NOW as u32), 8); let reply = ingest .handle(&datagram, peer(), NOW) .expect("a revoked token must be told, not ignored"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); let (_, msg) = otproto::open_message(&k_down, &reply).expect("sealed under K_down"); match msg { Message::Nack(nack) => assert_eq!(nack.reason, NackReason::UnknownToken), other => panic!("expected a NACK, got {other:?}"), } // And the points did not land. assert_eq!(ingest.counters.points_accepted.load(Ordering::Relaxed), 0); } /// The unauthenticated path, and everything that bounds it. #[tokio::test] async fn an_unknown_token_draws_one_small_rate_limited_notice() { let (ingest, _db, _dir, _task) = fixture().await; let stranger = 0xDEAD_BEEF_CAFE_F00D_u64; let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); let datagram = otproto::seal_message(&k_up, stranger, [9; 12], &loc(NOW as u32)); let reply = ingest .handle(&datagram, peer(), NOW) .expect("an unknown token should draw a notice"); // Never larger than what provoked it: that is what keeps a reflector from // being an amplifier. assert!( reply.len() <= datagram.len(), "reply {} B for a {} B request", reply.len(), datagram.len() ); // Sealed under this token id's K_rev, which no other device can derive. let k_rev = otproto::revocation_key(&REVOCATION_MASTER, stranger); let (header, msg) = otproto::open_message(&k_rev, &reply).expect("sealed under K_rev"); assert_eq!(header.token_id, stranger); assert_eq!( msg, Message::Revoked(Revoked { reason: RevokeReason::Unknown }) ); assert!( otproto::open( &otproto::revocation_key(&REVOCATION_MASTER, stranger ^ 1), &reply ) .is_err(), "a notice opened under another token's K_rev" ); // One per destination per minute. Everything after is silence, so a // spoofed victim is sent a message, not a flood. for i in 0..50 { assert!( ingest.handle(&datagram, peer(), NOW).is_none(), "a second notice went out on attempt {i}" ); } assert_eq!( ingest .counters .unverified_notices_sent .load(Ordering::Relaxed), 1 ); } /// A datagram too short to have cost the sender anything earns nothing. #[tokio::test] async fn a_minimum_size_datagram_never_draws_a_notice() { let (ingest, _db, _dir, _task) = fixture().await; let mut runt = vec![0u8; otproto::MIN_DATAGRAM]; runt[0] = 0x11; // version 1, type LOC runt[1..9].copy_from_slice(&0xDEAD_BEEF_u64.to_be_bytes()); assert!(ingest.handle(&runt, peer(), NOW).is_none()); assert_eq!( ingest .counters .unverified_notices_sent .load(Ordering::Relaxed), 0 ); } /// With no master configured, the server has no way to reply to something it /// cannot verify — which is the whole of the config switch. #[tokio::test] async fn notices_are_off_without_a_revocation_master() { let dir = tempfile::tempdir().expect("temp dir"); let db = Db::open(&dir.path().join("t.db")).await.expect("open"); let (writer, _task) = crate::writer::spawn(db.write.clone()); let ingest = Ingest::new(writer, 30 * 86_400, None); let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); let datagram = otproto::seal_message(&k_up, 0x1234, [4; 12], &loc(NOW as u32)); assert!(ingest.handle(&datagram, peer(), NOW).is_none()); } #[tokio::test] async fn a_config_pending_flag_appears_when_the_device_is_behind() { let (ingest, _db, _dir, _task) = fixture().await; ingest.insert_token(TokenSlot::new(TOKEN_ID, 1, &TOKEN_KEY, 5)); let hello = Message::Hello(otproto::Hello { app_version_code: 1, os_api_level: 34, flags: otproto::HelloFlags::NONE, config_version: 2, // behind the server's 5 }); let datagram = seal(&hello, 9); let reply = ingest.handle(&datagram, peer(), NOW).expect("ack"); let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); match otproto::open_message(&k_down, &reply).expect("opens").1 { Message::Ack(a) => assert!( a.flags.contains(AckFlags::CONFIG_PENDING), "the device is behind and must be told to fetch config" ), other => panic!("expected an ACK, got {other:?}"), } } /// The anti-amplification property, as it actually is. /// /// Not `reply <= request` — that rule was paid for with reserved padding on /// every `HELLO` and `PING`, to defend against a threat authentication already /// removes. What holds instead: every reply fits the reply budget, and the /// resulting ratio is nowhere near enough leverage to be worth reflecting /// through — and the sender needed a valid token key to get a reply at all, /// which the silence tests above cover. #[tokio::test] async fn replies_stay_inside_the_reply_budget() { let (ingest, _db, _dir, _task) = fixture().await; let requests = [ seal(&loc(NOW as u32), 10), seal( &Message::Loc( (0..MAX_POINTS as u32) .map(|i| Point::new(NOW as u32 - i, 1, 2)) .collect(), ), 11, ), seal(&Message::Ping(otproto::Ping { echo: 1, seq: 1 }), 12), seal( &Message::Hello(otproto::Hello { app_version_code: 1, os_api_level: 29, flags: otproto::HelloFlags::NONE, config_version: 1, }), 13, ), ]; for datagram in requests { if let Some(reply) = ingest.handle(&datagram, peer(), NOW) { assert!( otproto::fits_reply_budget(reply.len()), "a {}-byte request drew a {}-byte reply, over the {}-byte budget", datagram.len(), reply.len(), otproto::MAX_REPLY, ); let ratio = reply.len() as f64 / datagram.len() as f64; assert!( ratio <= 1.5, "a {}-byte request drew a {}-byte reply, {ratio:.2}x amplification", datagram.len(), reply.len(), ); } } } /// Garbage never panics, and never draws anything but a bounded notice. /// /// A run of `0x11` bytes parses as a well-formed header for token /// `0x1111111111111111`, which the server does not have — so the notice path /// is reachable from pure garbage by construction. That is expected. What /// must hold is that the reply is never larger than the request and that the /// budget stops it almost immediately. #[tokio::test] async fn garbage_never_panics_and_never_amplifies() { let (ingest, _db, _dir, _task) = fixture().await; for len in [0usize, 1, 20, 36, 37, 100, 1200, 1201] { for fill in [0u8, 0x11, 0xFF] { let datagram = vec![fill; len]; if let Some(reply) = ingest.handle(&datagram, peer(), NOW) { assert!( reply.len() <= datagram.len(), "len {len} fill {fill} drew a {} B reply", reply.len() ); } } } assert!( ingest .counters .unverified_notices_sent .load(Ordering::Relaxed) <= 1, "the per-destination budget should have stopped after the first notice" ); } }