ingest.rs
| 1 | //! Transport-agnostic datagram handling. |
| 2 | //! |
| 3 | //! [`Ingest::handle`] is **synchronous** and takes a couple of microseconds. It |
| 4 | //! never touches the database: every active token lives in an in-memory map, |
| 5 | //! loaded at startup and updated on mint/revoke. That is what lets the UDP |
| 6 | //! receive loop stay a tight `recv_from` → `handle` → `send_to` cycle with no |
| 7 | //! `.await` in the middle, and it is why an unknown `token_id` is genuinely |
| 8 | //! unknown rather than merely uncached. |
| 9 | //! |
| 10 | //! The cost is memory proportional to the number of active tokens. At ~90 bytes |
| 11 | //! per slot, a million tokens would be 90 MB; for a self-hosted instance with a |
| 12 | //! handful of users it is a few kilobytes. If that ever stops being true, this is |
| 13 | //! the module to revisit. |
| 14 | //! |
| 15 | //! The same function serves the UDP loop and the TLS fallback listener, so the |
| 16 | //! two transports cannot drift apart in their handling of anything. |
| 17 | |
| 18 | use std::net::{IpAddr, SocketAddr}; |
| 19 | use std::sync::Arc; |
| 20 | use std::sync::atomic::{AtomicU64, Ordering}; |
| 21 | |
| 22 | use dashmap::DashMap; |
| 23 | #[cfg(test)] |
| 24 | use otproto::msg::Direction; |
| 25 | use otproto::{ |
| 26 | Ack, AckFlags, DecodeError, Header, Key, MAX_POINTS, Message, Nack, NackReason, Point, |
| 27 | RevokeReason, Revoked, kdf, |
| 28 | }; |
| 29 | use rand::TryRngCore; |
| 30 | use rand::rngs::OsRng; |
| 31 | use tracing::{debug, trace}; |
| 32 | |
| 33 | use crate::limits::Limits; |
| 34 | use crate::writer::{Accepted, WriteHandle, WriteOp}; |
| 35 | |
| 36 | /// Which transport a datagram arrived on. Recorded per token so the UI can show |
| 37 | /// whether a phone is on UDP or has fallen back to TLS. |
| 38 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 39 | pub enum Transport { |
| 40 | Udp, |
| 41 | /// Constructed by the TLS-over-TCP fallback listener, which lands in a later |
| 42 | /// step; the recording path is already transport-agnostic so that listener |
| 43 | /// only has to call `Ingest::handle` with this variant. |
| 44 | #[allow(dead_code)] |
| 45 | Tls, |
| 46 | } |
| 47 | |
| 48 | impl Transport { |
| 49 | const fn as_str(self) -> &'static str { |
| 50 | match self { |
| 51 | Self::Udp => "udp", |
| 52 | Self::Tls => "tls", |
| 53 | } |
| 54 | } |
| 55 | } |
| 56 | |
| 57 | #[derive(Debug, Clone, Copy)] |
| 58 | pub struct Peer { |
| 59 | pub addr: SocketAddr, |
| 60 | pub transport: Transport, |
| 61 | } |
| 62 | |
| 63 | impl Peer { |
| 64 | pub fn ip(&self) -> IpAddr { |
| 65 | self.addr.ip() |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | /// Everything needed to serve one token, with no database round trip. |
| 70 | #[derive(Debug)] |
| 71 | pub struct TokenSlot { |
| 72 | pub token_id: u64, |
| 73 | /// Resolved once, at load time, so the write path never has to ask which |
| 74 | /// account a token belongs to. |
| 75 | pub user_id: i64, |
| 76 | pub k_up: Key, |
| 77 | pub k_down: Key, |
| 78 | pub config_version: u16, |
| 79 | /// Revoked tokens stay in the map instead of being removed. |
| 80 | /// |
| 81 | /// The key is what lets the server answer at all, and a revoked device is |
| 82 | /// precisely the one it needs to answer. Dropping the slot would leave the |
| 83 | /// only possible reply an unauthenticated one — a reflector — where keeping |
| 84 | /// it makes the reply a sealed `NACK` nobody else could have produced. |
| 85 | /// |
| 86 | /// The token authorises nothing from the moment this is set; the key survives |
| 87 | /// only so the server can say so. |
| 88 | pub revoked: Option<RevokeReason>, |
| 89 | } |
| 90 | |
| 91 | impl TokenSlot { |
| 92 | pub fn new(token_id: u64, user_id: i64, token_key: &Key, config_version: u16) -> Self { |
| 93 | let (k_up, k_down) = kdf::derive_both(token_key); |
| 94 | Self { |
| 95 | token_id, |
| 96 | user_id, |
| 97 | k_up, |
| 98 | k_down, |
| 99 | config_version, |
| 100 | revoked: None, |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | #[must_use] |
| 105 | pub fn revoked(mut self, reason: RevokeReason) -> Self { |
| 106 | self.revoked = Some(reason); |
| 107 | self |
| 108 | } |
| 109 | } |
| 110 | |
| 111 | /// Aggregate counters. Per-packet logging would itself be an amplifier, so these |
| 112 | /// are the only per-packet observability, exposed on a localhost-only `/metrics`. |
| 113 | #[derive(Debug, Default)] |
| 114 | pub struct Counters { |
| 115 | pub received: AtomicU64, |
| 116 | pub malformed: AtomicU64, |
| 117 | pub unknown_token: AtomicU64, |
| 118 | pub auth_failed: AtomicU64, |
| 119 | pub rate_limited: AtomicU64, |
| 120 | pub throttled: AtomicU64, |
| 121 | pub points_accepted: AtomicU64, |
| 122 | pub points_rejected: AtomicU64, |
| 123 | pub acks_sent: AtomicU64, |
| 124 | pub nacks_sent: AtomicU64, |
| 125 | pub silent_drops: AtomicU64, |
| 126 | /// Sealed notices to a token the server still holds a key for. |
| 127 | pub revoked_notices_sent: AtomicU64, |
| 128 | /// Notices sent without being able to verify the request. The reflection |
| 129 | /// budget, in other words — worth watching. |
| 130 | pub unverified_notices_sent: AtomicU64, |
| 131 | pub notices_suppressed: AtomicU64, |
| 132 | } |
| 133 | |
| 134 | impl Counters { |
| 135 | fn bump(counter: &AtomicU64) { |
| 136 | counter.fetch_add(1, Ordering::Relaxed); |
| 137 | } |
| 138 | |
| 139 | fn add(counter: &AtomicU64, n: u64) { |
| 140 | counter.fetch_add(n, Ordering::Relaxed); |
| 141 | } |
| 142 | } |
| 143 | |
| 144 | pub struct Ingest { |
| 145 | tokens: DashMap<u64, Arc<TokenSlot>>, |
| 146 | limits: Limits, |
| 147 | writer: WriteHandle, |
| 148 | pub counters: Counters, |
| 149 | /// Accept timestamps within ±this many seconds of the server clock. |
| 150 | ts_window_s: u32, |
| 151 | /// Master for per-token revocation keys, or `None` when unverified notices |
| 152 | /// are switched off in config. |
| 153 | /// |
| 154 | /// `Some` is the only thing that lets this server reply to a datagram it |
| 155 | /// cannot verify, so the config switch is represented as the presence of the |
| 156 | /// key rather than as a separate boolean. There is then no way to enable the |
| 157 | /// behaviour by accident. |
| 158 | revocation_master: Option<Key>, |
| 159 | } |
| 160 | |
| 161 | impl Ingest { |
| 162 | pub fn new(writer: WriteHandle, ts_window_s: u32, revocation_master: Option<Key>) -> Self { |
| 163 | Self { |
| 164 | tokens: DashMap::new(), |
| 165 | limits: Limits::new(), |
| 166 | writer, |
| 167 | counters: Counters::default(), |
| 168 | ts_window_s, |
| 169 | revocation_master, |
| 170 | } |
| 171 | } |
| 172 | |
| 173 | pub fn insert_token(&self, slot: TokenSlot) { |
| 174 | self.tokens.insert(slot.token_id, Arc::new(slot)); |
| 175 | } |
| 176 | |
| 177 | /// Called by revoke, by password change, and by the staleness sweep. |
| 178 | /// |
| 179 | /// The slot stays, flagged. The token stops authorising anything |
| 180 | /// immediately; what survives is only the ability to tell that one device it |
| 181 | /// is finished, in a message nobody else could have sealed. |
| 182 | pub fn mark_revoked(&self, token_id: u64, reason: RevokeReason) { |
| 183 | if let Some(existing) = self.tokens.get(&token_id).map(|s| Arc::clone(&s)) { |
| 184 | if existing.revoked.is_some() { |
| 185 | return; |
| 186 | } |
| 187 | self.tokens.insert( |
| 188 | token_id, |
| 189 | Arc::new(TokenSlot { |
| 190 | revoked: Some(reason), |
| 191 | ..*existing |
| 192 | }), |
| 193 | ); |
| 194 | } |
| 195 | } |
| 196 | |
| 197 | /// Active tokens only — a revoked slot is bookkeeping, not a live device. |
| 198 | pub fn active_token_count(&self) -> usize { |
| 199 | self.tokens.iter().filter(|s| s.revoked.is_none()).count() |
| 200 | } |
| 201 | |
| 202 | pub fn limits(&self) -> &Limits { |
| 203 | &self.limits |
| 204 | } |
| 205 | |
| 206 | /// Handle one datagram. Returns the bytes to send back, if any. |
| 207 | /// |
| 208 | /// `now` is passed in rather than read from the clock so this is testable |
| 209 | /// without sleeping. |
| 210 | pub fn handle(&self, datagram: &[u8], peer: Peer, now: i64) -> Option<Vec<u8>> { |
| 211 | Counters::bump(&self.counters.received); |
| 212 | |
| 213 | // 1. Structure. No state, no allocation, no crypto. |
| 214 | let header = match Header::peek(datagram) { |
| 215 | Ok(h) => h, |
| 216 | Err(e) => { |
| 217 | Counters::bump(&self.counters.malformed); |
| 218 | trace!(?e, "dropping malformed datagram"); |
| 219 | return self.silent(); |
| 220 | } |
| 221 | }; |
| 222 | |
| 223 | // A downlink type arriving on the uplink is either a bug or someone |
| 224 | // replaying our own traffic back at us. It can never be legitimate. |
| 225 | if !header.msg_type.is_uplink() { |
| 226 | Counters::bump(&self.counters.malformed); |
| 227 | return self.silent(); |
| 228 | } |
| 229 | |
| 230 | // 2. Per-IP budget and bans. |
| 231 | if let Err(reason) = self.limits.check_ip(peer.ip()) { |
| 232 | Counters::bump(&self.counters.rate_limited); |
| 233 | trace!(?reason, "dropping rate-limited datagram"); |
| 234 | return self.silent(); |
| 235 | } |
| 236 | |
| 237 | // 3. Does this token exist? |
| 238 | let Some(slot) = self.tokens.get(&header.token_id).map(|s| Arc::clone(&s)) else { |
| 239 | Counters::bump(&self.counters.unknown_token); |
| 240 | self.limits.note_unknown_token(peer.ip()); |
| 241 | // No key, so no way to verify this datagram — which makes any reply a |
| 242 | // reply to an address the sender merely claimed. See |
| 243 | // [`Self::unverified_notice`] for what makes that tolerable. |
| 244 | return self.unverified_notice(header.token_id, peer, datagram.len()); |
| 245 | }; |
| 246 | |
| 247 | // 4. AEAD. The first expensive step, ~1 µs for a 62-byte packet. |
| 248 | // |
| 249 | // Everything below this line may answer; nothing above it ever does. |
| 250 | // That is not a style rule, it is the anti-reflection defence. UDP source |
| 251 | // addresses are trivially forged, and `token_id` travels in cleartext, so |
| 252 | // anyone who has seen one datagram can name a valid token. If the server |
| 253 | // replied before verifying, an attacker could spoof a victim's address |
| 254 | // and have us send them a packet per junk datagram — laundering the |
| 255 | // attacker's origin and firing at whatever rate they choose. |
| 256 | // |
| 257 | // Which is why the per-token budget is checked *after* this and not |
| 258 | // before, even though that costs an AEAD open on every flooded packet. |
| 259 | // The per-IP budget above absorbs the bulk at no crypto cost. |
| 260 | let payload = match otproto::open(&slot.k_up, datagram) { |
| 261 | Ok((_, payload)) => payload, |
| 262 | Err(DecodeError::AuthFailed) => { |
| 263 | Counters::bump(&self.counters.auth_failed); |
| 264 | self.limits.note_aead_failure(peer.ip()); |
| 265 | // Never answer this. The server cannot know who sent it, so a |
| 266 | // reply would be both a forgery oracle and a reflector. |
| 267 | return self.silent(); |
| 268 | } |
| 269 | Err(_) => { |
| 270 | Counters::bump(&self.counters.malformed); |
| 271 | return self.silent(); |
| 272 | } |
| 273 | }; |
| 274 | |
| 275 | // 5. Is this token still alive? Checked after AEAD, so the answer is a |
| 276 | // sealed message the real device can trust and nobody else can forge. |
| 277 | // This is the whole reason revoked slots keep their keys. |
| 278 | if let Some(reason) = slot.revoked { |
| 279 | Counters::bump(&self.counters.revoked_notices_sent); |
| 280 | debug!(token_id = slot.token_id, ?reason, "revoked token reported"); |
| 281 | return self.seal_reply( |
| 282 | &slot, |
| 283 | Message::Nack(Nack { |
| 284 | nonce: header.nonce, |
| 285 | reason: NackReason::UnknownToken, |
| 286 | retry_after_s: 0, |
| 287 | }), |
| 288 | ); |
| 289 | } |
| 290 | |
| 291 | // 6. Per-token budget. Authenticated, so this NACK reaches the device |
| 292 | // that actually sent the datagram and nobody else. |
| 293 | if self.limits.check_token(header.token_id).is_err() { |
| 294 | Counters::bump(&self.counters.rate_limited); |
| 295 | return self.nack(&slot, header.nonce, NackReason::RateLimited, 5); |
| 296 | } |
| 297 | |
| 298 | let msg = match Message::decode_payload(header.msg_type, &payload) { |
| 299 | Ok(m) => m, |
| 300 | Err(e) => { |
| 301 | Counters::bump(&self.counters.malformed); |
| 302 | debug!( |
| 303 | token_id = header.token_id, |
| 304 | ?e, |
| 305 | "authenticated but malformed payload" |
| 306 | ); |
| 307 | return self.nack(&slot, header.nonce, NackReason::Malformed, 0); |
| 308 | } |
| 309 | }; |
| 310 | |
| 311 | // From here the packet is authenticated, so telemetry is safe to record. |
| 312 | self.record_seen(&slot, peer, now); |
| 313 | |
| 314 | match msg { |
| 315 | Message::Loc(points) => self.handle_loc(&slot, header, points, now), |
| 316 | Message::Hello(hello) => { |
| 317 | self.writer.try_send(WriteOp::TokenHello { |
| 318 | token_id: slot.token_id as i64, |
| 319 | app_version: i64::from(hello.app_version_code), |
| 320 | os_api_level: i64::from(hello.os_api_level), |
| 321 | }); |
| 322 | self.ack(&slot, vec![header.nonce], hello.config_version) |
| 323 | } |
| 324 | Message::ConfigGet(_) => { |
| 325 | // CONFIG is served by the HTTP/state path in this build; a device |
| 326 | // asking gets silence rather than a stale answer, and retries. |
| 327 | // Wired up with the config editor in a later step. |
| 328 | self.silent() |
| 329 | } |
| 330 | Message::Ping(ping) => { |
| 331 | // The echo is opaque to us; copying it back is the whole job. |
| 332 | let pong = Message::Pong(otproto::Pong { |
| 333 | echo: ping.echo, |
| 334 | seq: ping.seq, |
| 335 | }); |
| 336 | self.seal_reply(&slot, pong) |
| 337 | } |
| 338 | // Unreachable: the direction check above rejected every downlink type. |
| 339 | Message::Ack(_) |
| 340 | | Message::Nack(_) |
| 341 | | Message::Config(_) |
| 342 | | Message::Pong(_) |
| 343 | | Message::Revoked(_) => { |
| 344 | Counters::bump(&self.counters.malformed); |
| 345 | self.silent() |
| 346 | } |
| 347 | } |
| 348 | } |
| 349 | |
| 350 | fn handle_loc( |
| 351 | &self, |
| 352 | slot: &TokenSlot, |
| 353 | header: Header, |
| 354 | points: Vec<Point>, |
| 355 | now: i64, |
| 356 | ) -> Option<Vec<u8>> { |
| 357 | // Semantic validation. A point far outside the timestamp window would |
| 358 | // land where the retention sweep never reaches it, so it is dropped |
| 359 | // rather than stored — but the rest of the batch is still kept, because |
| 360 | // one bad fix must not cost the user a whole journey. |
| 361 | let before = points.len(); |
| 362 | let accepted: Vec<Point> = points |
| 363 | .into_iter() |
| 364 | .filter(|p| p.validate(now as u32, self.ts_window_s).is_ok()) |
| 365 | .collect(); |
| 366 | let rejected = before - accepted.len(); |
| 367 | Counters::add(&self.counters.points_rejected, rejected as u64); |
| 368 | |
| 369 | if accepted.is_empty() { |
| 370 | Counters::bump(&self.counters.malformed); |
| 371 | return self.nack(slot, header.nonce, NackReason::Malformed, 0); |
| 372 | } |
| 373 | Counters::add(&self.counters.points_accepted, accepted.len() as u64); |
| 374 | |
| 375 | let accepted_count = accepted.len(); |
| 376 | match self.writer.try_send(WriteOp::Points { |
| 377 | user_id: slot.user_id, |
| 378 | src_token_id: slot.token_id as i64, |
| 379 | points: accepted, |
| 380 | recv_at: now, |
| 381 | }) { |
| 382 | Accepted::Yes => {} |
| 383 | Accepted::Saturated => { |
| 384 | // Do not ack what we did not store: the client keeps the points |
| 385 | // queued and retries. THROTTLE tells it to slow down first. |
| 386 | Counters::bump(&self.counters.throttled); |
| 387 | return self.throttled(slot, header.nonce); |
| 388 | } |
| 389 | Accepted::Closed => return self.silent(), |
| 390 | } |
| 391 | |
| 392 | trace!( |
| 393 | token_id = slot.token_id, |
| 394 | user_id = slot.user_id, |
| 395 | accepted_count, |
| 396 | "stored points" |
| 397 | ); |
| 398 | self.ack(slot, vec![header.nonce], slot.config_version) |
| 399 | } |
| 400 | |
| 401 | fn record_seen(&self, slot: &TokenSlot, peer: Peer, now: i64) { |
| 402 | // Best effort: if the writer is saturated, telemetry is the first thing |
| 403 | // worth dropping. |
| 404 | self.writer.try_send(WriteOp::TokenSeen { |
| 405 | token_id: slot.token_id as i64, |
| 406 | at: now, |
| 407 | src_ip: peer.ip().to_string(), |
| 408 | src_port: peer.addr.port(), |
| 409 | transport: peer.transport.as_str(), |
| 410 | }); |
| 411 | } |
| 412 | |
| 413 | fn ack( |
| 414 | &self, |
| 415 | slot: &TokenSlot, |
| 416 | nonces: Vec<[u8; 12]>, |
| 417 | device_config_version: u16, |
| 418 | ) -> Option<Vec<u8>> { |
| 419 | let flags = if device_config_version < slot.config_version { |
| 420 | AckFlags::CONFIG_PENDING |
| 421 | } else { |
| 422 | AckFlags::NONE |
| 423 | }; |
| 424 | debug_assert!(nonces.len() <= MAX_POINTS); |
| 425 | let ack = Message::Ack(Ack { nonces, flags }); |
| 426 | Counters::bump(&self.counters.acks_sent); |
| 427 | self.seal_reply(slot, ack) |
| 428 | } |
| 429 | |
| 430 | /// The writer channel is full, so the points were not stored. |
| 431 | /// |
| 432 | /// This is a `NACK`, not an `ACK` with the `THROTTLE` flag, because an `ACK` |
| 433 | /// retires the nonces it names: acking here would tell the client to delete |
| 434 | /// points that never reached the database. `RateLimited` with a retry hint is |
| 435 | /// exactly the "keep them and back off" semantic the client already |
| 436 | /// implements, so saturation costs a delay rather than data. |
| 437 | fn throttled(&self, slot: &TokenSlot, nonce: [u8; 12]) -> Option<Vec<u8>> { |
| 438 | self.nack(slot, nonce, NackReason::RateLimited, 10) |
| 439 | } |
| 440 | |
| 441 | fn nack( |
| 442 | &self, |
| 443 | slot: &TokenSlot, |
| 444 | nonce: [u8; 12], |
| 445 | reason: NackReason, |
| 446 | retry_after_s: u8, |
| 447 | ) -> Option<Vec<u8>> { |
| 448 | Counters::bump(&self.counters.nacks_sent); |
| 449 | self.seal_reply( |
| 450 | slot, |
| 451 | Message::Nack(Nack { |
| 452 | nonce, |
| 453 | reason, |
| 454 | retry_after_s, |
| 455 | }), |
| 456 | ) |
| 457 | } |
| 458 | |
| 459 | /// The one reply this server sends without having verified the request. |
| 460 | /// |
| 461 | /// The server has no record of this `token_id`, so it cannot open the |
| 462 | /// datagram and cannot know who really sent it. `peer` is whatever the source |
| 463 | /// address claimed, which on UDP is forgeable. Answering therefore makes this |
| 464 | /// port a reflector, and three things are what keep that from mattering: |
| 465 | /// |
| 466 | /// 1. **It cannot amplify.** [`Revoked`] is 38 bytes, the smallest message in |
| 467 | /// the protocol, and a request shorter than that gets nothing. |
| 468 | /// 2. **It is rate limited by destination.** The budget is keyed on the |
| 469 | /// address the reply would go to — the victim, for a spoofed packet — at |
| 470 | /// one per minute, under a global ceiling for distributed attempts. |
| 471 | /// 3. **It is optional.** No revocation master configured, no reply. |
| 472 | /// |
| 473 | /// It cannot be forged, which is the other half. `K_rev` is derived from a |
| 474 | /// server master and this `token_id`, so a third party cannot produce one and |
| 475 | /// neither can another legitimate device — each only ever learns its own. |
| 476 | /// |
| 477 | /// This path exists for the cases where the row is genuinely gone: a database |
| 478 | /// restored from a backup that predates the login, or a rotated server key. |
| 479 | /// The ordinary revocation path keeps the slot and answers with a sealed |
| 480 | /// `NACK`, which needs none of the above. |
| 481 | fn unverified_notice(&self, token_id: u64, peer: Peer, request_len: usize) -> Option<Vec<u8>> { |
| 482 | let Some(master) = self.revocation_master else { |
| 483 | return self.silent(); |
| 484 | }; |
| 485 | if !otproto::may_answer_unverified(request_len) { |
| 486 | return self.silent(); |
| 487 | } |
| 488 | if self.limits.check_notice(peer.ip()).is_err() { |
| 489 | Counters::bump(&self.counters.notices_suppressed); |
| 490 | return self.silent(); |
| 491 | } |
| 492 | |
| 493 | let mut nonce = [0u8; 12]; |
| 494 | if OsRng.try_fill_bytes(&mut nonce).is_err() { |
| 495 | return self.silent(); |
| 496 | } |
| 497 | let msg = Message::Revoked(Revoked { |
| 498 | reason: RevokeReason::Unknown, |
| 499 | }); |
| 500 | let reply = otproto::seal_message( |
| 501 | &otproto::revocation_key(&master, token_id), |
| 502 | token_id, |
| 503 | nonce, |
| 504 | &msg, |
| 505 | ); |
| 506 | debug_assert!( |
| 507 | reply.len() <= request_len, |
| 508 | "an unverified reply must never exceed its request" |
| 509 | ); |
| 510 | Counters::bump(&self.counters.unverified_notices_sent); |
| 511 | Some(reply) |
| 512 | } |
| 513 | |
| 514 | /// Seal a reply. |
| 515 | /// |
| 516 | /// Anti-amplification is not enforced here, because by this point it is |
| 517 | /// already established: a reply is only reached after the datagram passed AEAD |
| 518 | /// verification, so a reflection attacker must hold a live token key, and |
| 519 | /// every reply this server can construct fits in [`otproto::MAX_REPLY`] bytes |
| 520 | /// — a ratio near 1 against any request. The `debug_assert` is a development |
| 521 | /// guard against a future message type outgrowing that budget, not a runtime |
| 522 | /// control. |
| 523 | fn seal_reply(&self, slot: &TokenSlot, msg: Message) -> Option<Vec<u8>> { |
| 524 | debug_assert!( |
| 525 | otproto::fits_reply_budget(otproto::datagram_len(msg.payload_len())), |
| 526 | "{:?} exceeds the reply budget of {} bytes", |
| 527 | msg.msg_type(), |
| 528 | otproto::MAX_REPLY, |
| 529 | ); |
| 530 | |
| 531 | let mut nonce = [0u8; 12]; |
| 532 | if OsRng.try_fill_bytes(&mut nonce).is_err() { |
| 533 | // Without fresh randomness we must not seal anything: reusing a nonce |
| 534 | // under ChaCha20-Poly1305 leaks the keystream. |
| 535 | return self.silent(); |
| 536 | } |
| 537 | Some(otproto::seal_message( |
| 538 | &slot.k_down, |
| 539 | slot.token_id, |
| 540 | nonce, |
| 541 | &msg, |
| 542 | )) |
| 543 | } |
| 544 | |
| 545 | fn silent(&self) -> Option<Vec<u8>> { |
| 546 | Counters::bump(&self.counters.silent_drops); |
| 547 | None |
| 548 | } |
| 549 | |
| 550 | /// Derived key for a direction. Test-only: the live paths look up the slot and |
| 551 | /// use both keys, and exposing a key by id elsewhere would invite misuse. |
| 552 | #[cfg(test)] |
| 553 | pub fn key_for(&self, token_id: u64, dir: Direction) -> Option<Key> { |
| 554 | self.tokens.get(&token_id).map(|s| match dir { |
| 555 | Direction::Up => s.k_up, |
| 556 | Direction::Down => s.k_down, |
| 557 | }) |
| 558 | } |
| 559 | } |
| 560 | |
| 561 | #[cfg(test)] |
| 562 | mod tests { |
| 563 | use super::*; |
| 564 | use crate::db::Db; |
| 565 | use otproto::point::Flags; |
| 566 | use otproto::{HEADER_LEN, TAG_LEN}; |
| 567 | |
| 568 | const TOKEN_ID: u64 = 0x0123_4567_89AB_CDEF; |
| 569 | const TOKEN_KEY: Key = [0x5A; 32]; |
| 570 | const REVOCATION_MASTER: Key = [0xC3; 32]; |
| 571 | const NOW: i64 = 1_785_000_042; |
| 572 | |
| 573 | fn peer() -> Peer { |
| 574 | Peer { |
| 575 | addr: "203.0.113.5:40000".parse().expect("literal"), |
| 576 | transport: Transport::Udp, |
| 577 | } |
| 578 | } |
| 579 | |
| 580 | async fn fixture() -> (Ingest, Db, tempfile::TempDir, tokio::task::JoinHandle<()>) { |
| 581 | let dir = tempfile::tempdir().expect("temp dir"); |
| 582 | let db = Db::open(&dir.path().join("t.db")).await.expect("open"); |
| 583 | sqlx::query( |
| 584 | "INSERT INTO users (id, username, pw_hash, display_name, created_at, pw_changed_at) \ |
| 585 | VALUES (1, 'a', 'x', 'A', 0, 0)", |
| 586 | ) |
| 587 | .execute(&db.write) |
| 588 | .await |
| 589 | .expect("user"); |
| 590 | |
| 591 | let (writer, task) = crate::writer::spawn(db.write.clone()); |
| 592 | let ingest = Ingest::new(writer, 30 * 86_400, Some(REVOCATION_MASTER)); |
| 593 | ingest.insert_token(TokenSlot::new(TOKEN_ID, 1, &TOKEN_KEY, 1)); |
| 594 | (ingest, db, dir, task) |
| 595 | } |
| 596 | |
| 597 | fn seal(msg: &Message, nonce_seed: u8) -> Vec<u8> { |
| 598 | let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); |
| 599 | otproto::seal_message(&k_up, TOKEN_ID, [nonce_seed; 12], msg) |
| 600 | } |
| 601 | |
| 602 | fn loc(ts: u32) -> Message { |
| 603 | Message::Loc(vec![Point { |
| 604 | acc_dm: Some(80), |
| 605 | flags: Flags::NONE, |
| 606 | ..Point::new(ts, 525_200_080, 134_050_000) |
| 607 | }]) |
| 608 | } |
| 609 | |
| 610 | /// The anti-reflection invariant, stated as a test: nothing that fails the |
| 611 | /// AEAD check may produce a reply. |
| 612 | /// |
| 613 | /// UDP source addresses are forgeable, so any reply to an unverified |
| 614 | /// datagram is a packet an attacker can aim at a third party. `token_id` is |
| 615 | /// cleartext, so naming a real token costs nothing — which is what makes the |
| 616 | /// per-token rate limit the interesting case here rather than a theoretical |
| 617 | /// one. Its budget is small enough that a flood trips it immediately. |
| 618 | #[tokio::test] |
| 619 | async fn nothing_that_fails_aead_ever_gets_a_reply() { |
| 620 | let (ingest, _db, _dir, _task) = fixture().await; |
| 621 | let valid = seal(&loc(NOW as u32), 1); |
| 622 | |
| 623 | // Far past any per-token budget, so the pre-AEAD ordering bug would show |
| 624 | // up here as a rate-limit NACK sent to an unauthenticated sender. |
| 625 | for i in 0..200 { |
| 626 | // A real header naming a real token, with a corrupted tag. |
| 627 | let mut forged = valid.clone(); |
| 628 | let last = forged.len() - 1; |
| 629 | forged[last] ^= 1; |
| 630 | forged[HEADER_LEN] ^= i as u8; |
| 631 | assert_eq!( |
| 632 | ingest.handle(&forged, peer(), NOW), |
| 633 | None, |
| 634 | "a datagram that fails AEAD was answered on attempt {i}" |
| 635 | ); |
| 636 | |
| 637 | // And a header-only datagram, which cannot authenticate at all. |
| 638 | let stub = valid[..HEADER_LEN + TAG_LEN].to_vec(); |
| 639 | assert_eq!( |
| 640 | ingest.handle(&stub, peer(), NOW), |
| 641 | None, |
| 642 | "a truncated datagram was answered on attempt {i}" |
| 643 | ); |
| 644 | } |
| 645 | |
| 646 | assert_eq!( |
| 647 | ingest.counters.acks_sent.load(Ordering::Relaxed) |
| 648 | + ingest.counters.nacks_sent.load(Ordering::Relaxed), |
| 649 | 0, |
| 650 | "the server sent something in response to unauthenticated traffic" |
| 651 | ); |
| 652 | } |
| 653 | |
| 654 | #[tokio::test] |
| 655 | async fn a_valid_loc_is_acked_and_stored() { |
| 656 | let (ingest, db, _dir, _task) = fixture().await; |
| 657 | let datagram = seal(&loc(NOW as u32), 1); |
| 658 | let reply = ingest.handle(&datagram, peer(), NOW).expect("should ack"); |
| 659 | |
| 660 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 661 | let (_, msg) = otproto::open_message(&k_down, &reply).expect("ack opens"); |
| 662 | match msg { |
| 663 | Message::Ack(ack) => { |
| 664 | assert_eq!( |
| 665 | ack.nonces, |
| 666 | vec![[1u8; 12]], |
| 667 | "the ack must echo the request nonce" |
| 668 | ); |
| 669 | } |
| 670 | other => panic!("expected an ACK, got {other:?}"), |
| 671 | } |
| 672 | |
| 673 | tokio::time::sleep(std::time::Duration::from_millis(400)).await; |
| 674 | let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points WHERE user_id = 1") |
| 675 | .fetch_one(&db.read) |
| 676 | .await |
| 677 | .expect("count"); |
| 678 | assert_eq!(count, 1); |
| 679 | } |
| 680 | |
| 681 | #[tokio::test] |
| 682 | async fn an_unknown_token_is_counted_and_struck() { |
| 683 | let (ingest, _db, _dir, _task) = fixture().await; |
| 684 | let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); |
| 685 | let datagram = otproto::seal_message(&k_up, 0xDEAD_BEEF, [2; 12], &loc(NOW as u32)); |
| 686 | // The notice itself is covered by |
| 687 | // `an_unknown_token_draws_one_small_rate_limited_notice`; what matters |
| 688 | // here is that answering did not stop the sender being treated as a |
| 689 | // scanner. Naming unknown ids still earns strikes and eventually a ban. |
| 690 | let _ = ingest.handle(&datagram, peer(), NOW); |
| 691 | assert_eq!(ingest.counters.unknown_token.load(Ordering::Relaxed), 1); |
| 692 | assert_eq!(ingest.counters.points_accepted.load(Ordering::Relaxed), 0); |
| 693 | } |
| 694 | |
| 695 | #[tokio::test] |
| 696 | async fn a_forged_datagram_gets_silence() { |
| 697 | let (ingest, _db, _dir, _task) = fixture().await; |
| 698 | let mut datagram = seal(&loc(NOW as u32), 3); |
| 699 | let last = datagram.len() - 1; |
| 700 | datagram[last] ^= 1; |
| 701 | assert!( |
| 702 | ingest.handle(&datagram, peer(), NOW).is_none(), |
| 703 | "a failed tag must never be answered" |
| 704 | ); |
| 705 | assert_eq!(ingest.counters.auth_failed.load(Ordering::Relaxed), 1); |
| 706 | } |
| 707 | |
| 708 | #[tokio::test] |
| 709 | async fn a_downlink_message_arriving_on_the_uplink_is_dropped() { |
| 710 | let (ingest, _db, _dir, _task) = fixture().await; |
| 711 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 712 | let pong = Message::Pong(otproto::Pong { echo: 1, seq: 3 }); |
| 713 | let datagram = otproto::seal_message(&k_down, TOKEN_ID, [4; 12], &pong); |
| 714 | assert!(ingest.handle(&datagram, peer(), NOW).is_none()); |
| 715 | } |
| 716 | |
| 717 | #[tokio::test] |
| 718 | async fn a_ping_is_answered_with_a_pong_of_no_greater_size() { |
| 719 | let (ingest, _db, _dir, _task) = fixture().await; |
| 720 | let ping = Message::Ping(otproto::Ping { |
| 721 | echo: 0xDEAD_BEEF, |
| 722 | seq: 7, |
| 723 | }); |
| 724 | let datagram = seal(&ping, 5); |
| 725 | let reply = ingest.handle(&datagram, peer(), NOW).expect("pong"); |
| 726 | assert!(reply.len() <= datagram.len(), "PONG amplified the PING"); |
| 727 | |
| 728 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 729 | match otproto::open_message(&k_down, &reply).expect("opens").1 { |
| 730 | Message::Pong(p) => { |
| 731 | assert_eq!(p.seq, 7); |
| 732 | assert_eq!( |
| 733 | p.echo, 0xDEAD_BEEF, |
| 734 | "the PING's opaque echo must come back untouched" |
| 735 | ); |
| 736 | } |
| 737 | other => panic!("expected a PONG, got {other:?}"), |
| 738 | } |
| 739 | } |
| 740 | |
| 741 | #[tokio::test] |
| 742 | async fn a_point_outside_the_timestamp_window_is_rejected() { |
| 743 | let (ingest, db, _dir, _task) = fixture().await; |
| 744 | // Year 2100: far beyond anything the retention sweep would ever reach. |
| 745 | let datagram = seal(&loc(4_102_444_800), 6); |
| 746 | let reply = ingest.handle(&datagram, peer(), NOW).expect("nack"); |
| 747 | |
| 748 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 749 | match otproto::open_message(&k_down, &reply).expect("opens").1 { |
| 750 | Message::Nack(n) => assert_eq!(n.reason, NackReason::Malformed), |
| 751 | other => panic!("expected a NACK, got {other:?}"), |
| 752 | } |
| 753 | |
| 754 | tokio::time::sleep(std::time::Duration::from_millis(300)).await; |
| 755 | let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") |
| 756 | .fetch_one(&db.read) |
| 757 | .await |
| 758 | .expect("count"); |
| 759 | assert_eq!(count, 0); |
| 760 | } |
| 761 | |
| 762 | #[tokio::test] |
| 763 | async fn one_bad_point_does_not_cost_the_whole_batch() { |
| 764 | let (ingest, db, _dir, _task) = fixture().await; |
| 765 | let good = Point { |
| 766 | acc_dm: Some(80), |
| 767 | ..Point::new(NOW as u32, 525_200_080, 134_050_000) |
| 768 | }; |
| 769 | let bad = Point::new(NOW as u32, 910_000_000, 0); // impossible latitude |
| 770 | let datagram = seal(&Message::Loc(vec![good, bad]), 7); |
| 771 | assert!( |
| 772 | ingest.handle(&datagram, peer(), NOW).is_some(), |
| 773 | "should still ack" |
| 774 | ); |
| 775 | |
| 776 | tokio::time::sleep(std::time::Duration::from_millis(400)).await; |
| 777 | let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM points") |
| 778 | .fetch_one(&db.read) |
| 779 | .await |
| 780 | .expect("count"); |
| 781 | assert_eq!(count, 1, "the good point must survive its bad neighbour"); |
| 782 | assert_eq!(ingest.counters.points_rejected.load(Ordering::Relaxed), 1); |
| 783 | } |
| 784 | |
| 785 | /// A revoked token keeps its key so this answer is possible at all. |
| 786 | /// |
| 787 | /// The alternative — dropping the slot — leaves silence as the only safe |
| 788 | /// reply, and the phone keeps reporting into nothing until someone notices. |
| 789 | #[tokio::test] |
| 790 | async fn a_revoked_token_is_told_so_in_a_message_it_can_verify() { |
| 791 | let (ingest, _db, _dir, _task) = fixture().await; |
| 792 | ingest.mark_revoked(TOKEN_ID, RevokeReason::Revoked); |
| 793 | |
| 794 | let datagram = seal(&loc(NOW as u32), 8); |
| 795 | let reply = ingest |
| 796 | .handle(&datagram, peer(), NOW) |
| 797 | .expect("a revoked token must be told, not ignored"); |
| 798 | |
| 799 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 800 | let (_, msg) = otproto::open_message(&k_down, &reply).expect("sealed under K_down"); |
| 801 | match msg { |
| 802 | Message::Nack(nack) => assert_eq!(nack.reason, NackReason::UnknownToken), |
| 803 | other => panic!("expected a NACK, got {other:?}"), |
| 804 | } |
| 805 | |
| 806 | // And the points did not land. |
| 807 | assert_eq!(ingest.counters.points_accepted.load(Ordering::Relaxed), 0); |
| 808 | } |
| 809 | |
| 810 | /// The unauthenticated path, and everything that bounds it. |
| 811 | #[tokio::test] |
| 812 | async fn an_unknown_token_draws_one_small_rate_limited_notice() { |
| 813 | let (ingest, _db, _dir, _task) = fixture().await; |
| 814 | let stranger = 0xDEAD_BEEF_CAFE_F00D_u64; |
| 815 | let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); |
| 816 | let datagram = otproto::seal_message(&k_up, stranger, [9; 12], &loc(NOW as u32)); |
| 817 | |
| 818 | let reply = ingest |
| 819 | .handle(&datagram, peer(), NOW) |
| 820 | .expect("an unknown token should draw a notice"); |
| 821 | |
| 822 | // Never larger than what provoked it: that is what keeps a reflector from |
| 823 | // being an amplifier. |
| 824 | assert!( |
| 825 | reply.len() <= datagram.len(), |
| 826 | "reply {} B for a {} B request", |
| 827 | reply.len(), |
| 828 | datagram.len() |
| 829 | ); |
| 830 | |
| 831 | // Sealed under this token id's K_rev, which no other device can derive. |
| 832 | let k_rev = otproto::revocation_key(&REVOCATION_MASTER, stranger); |
| 833 | let (header, msg) = otproto::open_message(&k_rev, &reply).expect("sealed under K_rev"); |
| 834 | assert_eq!(header.token_id, stranger); |
| 835 | assert_eq!( |
| 836 | msg, |
| 837 | Message::Revoked(Revoked { |
| 838 | reason: RevokeReason::Unknown |
| 839 | }) |
| 840 | ); |
| 841 | assert!( |
| 842 | otproto::open( |
| 843 | &otproto::revocation_key(&REVOCATION_MASTER, stranger ^ 1), |
| 844 | &reply |
| 845 | ) |
| 846 | .is_err(), |
| 847 | "a notice opened under another token's K_rev" |
| 848 | ); |
| 849 | |
| 850 | // One per destination per minute. Everything after is silence, so a |
| 851 | // spoofed victim is sent a message, not a flood. |
| 852 | for i in 0..50 { |
| 853 | assert!( |
| 854 | ingest.handle(&datagram, peer(), NOW).is_none(), |
| 855 | "a second notice went out on attempt {i}" |
| 856 | ); |
| 857 | } |
| 858 | assert_eq!( |
| 859 | ingest |
| 860 | .counters |
| 861 | .unverified_notices_sent |
| 862 | .load(Ordering::Relaxed), |
| 863 | 1 |
| 864 | ); |
| 865 | } |
| 866 | |
| 867 | /// A datagram too short to have cost the sender anything earns nothing. |
| 868 | #[tokio::test] |
| 869 | async fn a_minimum_size_datagram_never_draws_a_notice() { |
| 870 | let (ingest, _db, _dir, _task) = fixture().await; |
| 871 | let mut runt = vec![0u8; otproto::MIN_DATAGRAM]; |
| 872 | runt[0] = 0x11; // version 1, type LOC |
| 873 | runt[1..9].copy_from_slice(&0xDEAD_BEEF_u64.to_be_bytes()); |
| 874 | assert!(ingest.handle(&runt, peer(), NOW).is_none()); |
| 875 | assert_eq!( |
| 876 | ingest |
| 877 | .counters |
| 878 | .unverified_notices_sent |
| 879 | .load(Ordering::Relaxed), |
| 880 | 0 |
| 881 | ); |
| 882 | } |
| 883 | |
| 884 | /// With no master configured, the server has no way to reply to something it |
| 885 | /// cannot verify — which is the whole of the config switch. |
| 886 | #[tokio::test] |
| 887 | async fn notices_are_off_without_a_revocation_master() { |
| 888 | let dir = tempfile::tempdir().expect("temp dir"); |
| 889 | let db = Db::open(&dir.path().join("t.db")).await.expect("open"); |
| 890 | let (writer, _task) = crate::writer::spawn(db.write.clone()); |
| 891 | let ingest = Ingest::new(writer, 30 * 86_400, None); |
| 892 | |
| 893 | let k_up = kdf::derive(&TOKEN_KEY, Direction::Up); |
| 894 | let datagram = otproto::seal_message(&k_up, 0x1234, [4; 12], &loc(NOW as u32)); |
| 895 | assert!(ingest.handle(&datagram, peer(), NOW).is_none()); |
| 896 | } |
| 897 | |
| 898 | #[tokio::test] |
| 899 | async fn a_config_pending_flag_appears_when_the_device_is_behind() { |
| 900 | let (ingest, _db, _dir, _task) = fixture().await; |
| 901 | ingest.insert_token(TokenSlot::new(TOKEN_ID, 1, &TOKEN_KEY, 5)); |
| 902 | |
| 903 | let hello = Message::Hello(otproto::Hello { |
| 904 | app_version_code: 1, |
| 905 | os_api_level: 34, |
| 906 | flags: otproto::HelloFlags::NONE, |
| 907 | config_version: 2, // behind the server's 5 |
| 908 | }); |
| 909 | let datagram = seal(&hello, 9); |
| 910 | let reply = ingest.handle(&datagram, peer(), NOW).expect("ack"); |
| 911 | let k_down = kdf::derive(&TOKEN_KEY, Direction::Down); |
| 912 | match otproto::open_message(&k_down, &reply).expect("opens").1 { |
| 913 | Message::Ack(a) => assert!( |
| 914 | a.flags.contains(AckFlags::CONFIG_PENDING), |
| 915 | "the device is behind and must be told to fetch config" |
| 916 | ), |
| 917 | other => panic!("expected an ACK, got {other:?}"), |
| 918 | } |
| 919 | } |
| 920 | |
| 921 | /// The anti-amplification property, as it actually is. |
| 922 | /// |
| 923 | /// Not `reply <= request` — that rule was paid for with reserved padding on |
| 924 | /// every `HELLO` and `PING`, to defend against a threat authentication already |
| 925 | /// removes. What holds instead: every reply fits the reply budget, and the |
| 926 | /// resulting ratio is nowhere near enough leverage to be worth reflecting |
| 927 | /// through — and the sender needed a valid token key to get a reply at all, |
| 928 | /// which the silence tests above cover. |
| 929 | #[tokio::test] |
| 930 | async fn replies_stay_inside_the_reply_budget() { |
| 931 | let (ingest, _db, _dir, _task) = fixture().await; |
| 932 | let requests = [ |
| 933 | seal(&loc(NOW as u32), 10), |
| 934 | seal( |
| 935 | &Message::Loc( |
| 936 | (0..MAX_POINTS as u32) |
| 937 | .map(|i| Point::new(NOW as u32 - i, 1, 2)) |
| 938 | .collect(), |
| 939 | ), |
| 940 | 11, |
| 941 | ), |
| 942 | seal(&Message::Ping(otproto::Ping { echo: 1, seq: 1 }), 12), |
| 943 | seal( |
| 944 | &Message::Hello(otproto::Hello { |
| 945 | app_version_code: 1, |
| 946 | os_api_level: 29, |
| 947 | flags: otproto::HelloFlags::NONE, |
| 948 | config_version: 1, |
| 949 | }), |
| 950 | 13, |
| 951 | ), |
| 952 | ]; |
| 953 | for datagram in requests { |
| 954 | if let Some(reply) = ingest.handle(&datagram, peer(), NOW) { |
| 955 | assert!( |
| 956 | otproto::fits_reply_budget(reply.len()), |
| 957 | "a {}-byte request drew a {}-byte reply, over the {}-byte budget", |
| 958 | datagram.len(), |
| 959 | reply.len(), |
| 960 | otproto::MAX_REPLY, |
| 961 | ); |
| 962 | let ratio = reply.len() as f64 / datagram.len() as f64; |
| 963 | assert!( |
| 964 | ratio <= 1.5, |
| 965 | "a {}-byte request drew a {}-byte reply, {ratio:.2}x amplification", |
| 966 | datagram.len(), |
| 967 | reply.len(), |
| 968 | ); |
| 969 | } |
| 970 | } |
| 971 | } |
| 972 | |
| 973 | /// Garbage never panics, and never draws anything but a bounded notice. |
| 974 | /// |
| 975 | /// A run of `0x11` bytes parses as a well-formed header for token |
| 976 | /// `0x1111111111111111`, which the server does not have — so the notice path |
| 977 | /// is reachable from pure garbage by construction. That is expected. What |
| 978 | /// must hold is that the reply is never larger than the request and that the |
| 979 | /// budget stops it almost immediately. |
| 980 | #[tokio::test] |
| 981 | async fn garbage_never_panics_and_never_amplifies() { |
| 982 | let (ingest, _db, _dir, _task) = fixture().await; |
| 983 | for len in [0usize, 1, 20, 36, 37, 100, 1200, 1201] { |
| 984 | for fill in [0u8, 0x11, 0xFF] { |
| 985 | let datagram = vec![fill; len]; |
| 986 | if let Some(reply) = ingest.handle(&datagram, peer(), NOW) { |
| 987 | assert!( |
| 988 | reply.len() <= datagram.len(), |
| 989 | "len {len} fill {fill} drew a {} B reply", |
| 990 | reply.len() |
| 991 | ); |
| 992 | } |
| 993 | } |
| 994 | } |
| 995 | assert!( |
| 996 | ingest |
| 997 | .counters |
| 998 | .unverified_notices_sent |
| 999 | .load(Ordering::Relaxed) |
| 1000 | <= 1, |
| 1001 | "the per-destination budget should have stopped after the first notice" |
| 1002 | ); |
| 1003 | } |
| 1004 | } |
| 1005 |