ingest.rs
⎇
Raw
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
18use std::net::{IpAddr, SocketAddr};
19use std::sync::Arc;
20use std::sync::atomic::{AtomicU64, Ordering};
21
22use dashmap::DashMap;
23#[cfg(test)]
24use otproto::msg::Direction;
25use otproto::{
26 Ack, AckFlags, DecodeError, Header, Key, MAX_POINTS, Message, Nack, NackReason, Point,
27 RevokeReason, Revoked, kdf,
28};
29use rand::TryRngCore;
30use rand::rngs::OsRng;
31use tracing::{debug, trace};
32
33use crate::limits::Limits;
34use 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)]
39pub 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
48impl 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)]
58pub struct Peer {
59 pub addr: SocketAddr,
60 pub transport: Transport,
61}
62
63impl 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)]
71pub 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
91impl 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)]
114pub 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
134impl 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
144pub 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
161impl 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)]
562mod 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