//! `--simulate-device`: an in-process fake phone. //! //! It logs in over the real HTTP API, then walks a synthetic route sending real //! OTP/1 datagrams over a real UDP socket on loopback. Nothing is stubbed: the //! same codec, the same AEAD, the same ingest path, the same writer. //! //! This is the fastest feedback loop for everything downstream of the protocol — //! the web UI can be built against a moving marker before a line of Kotlin runs — //! and because it exercises the genuine encoder, a protocol mistake shows up here //! rather than on a phone. use std::time::Duration; use anyhow::{Context, Result, bail}; use otproto::msg::Direction; use otproto::point::Flags; use otproto::{AckFlags, Key, Message, Point, kdf}; use rand::TryRngCore; use rand::rngs::OsRng; use tokio::net::UdpSocket; use tracing::{info, warn}; /// How often the simulated phone reports. Matches the Balanced profile's walking /// cadence closely enough to be representative. const REPORT_INTERVAL: Duration = Duration::from_secs(5); /// Roughly walking pace, in units of 1e-7 degrees per report. const STEP_E7: i32 = 1_200; pub struct SimulatedDevice { pub base_url: String, pub username: String, pub password: String, pub udp_addr: String, } struct Credentials { token_id: u64, k_up: Key, k_down: Key, } impl SimulatedDevice { pub fn spawn(self) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { if let Err(e) = self.run().await { warn!(error = %e, "simulated device stopped"); } }) } async fn run(self) -> Result<()> { let creds = self.login().await?; let socket = UdpSocket::bind("0.0.0.0:0") .await .context("binding a client socket")?; socket .connect(&self.udp_addr) .await .with_context(|| format!("connecting to {}", self.udp_addr))?; info!(token_id = creds.token_id, addr = %self.udp_addr, "simulated device connected"); // Say hello once, exactly as the app does, so the server has a version // recorded against the token. let hello = Message::Hello(otproto::Hello { app_version_code: 0, os_api_level: 0, flags: otproto::HelloFlags::FIRST_LAUNCH, config_version: 1, }); self.send(&socket, &creds, &hello).await?; // A lap around a small block near the Brandenburg Gate. let mut lat = 525_200_080i32; let mut lon = 134_050_000i32; let mut leg = 0u32; let mut ticker = tokio::time::interval(REPORT_INTERVAL); loop { ticker.tick().await; // Four legs of 20 reports each: north, east, south, west. let (dlat, dlon) = match (leg / 20) % 4 { 0 => (STEP_E7, 0), 1 => (0, STEP_E7), 2 => (-STEP_E7, 0), _ => (0, -STEP_E7), }; lat += dlat; lon += dlon; leg = leg.wrapping_add(1); let bearing = match (leg / 20) % 4 { 0 => 0u16, 1 => 9_000, 2 => 18_000, _ => 27_000, }; let point = Point { acc_dm: Some(60), alt_m: Some(34), spd_cms: Some(140), brg_cdeg: Some(bearing), bat_pct: Some(80u8.saturating_sub((leg / 60) as u8)), flags: Flags::NONE, ..Point::new(crate::db::now() as u32, lat, lon) }; self.send(&socket, &creds, &Message::Loc(vec![point])) .await?; } } async fn login(&self) -> Result { let client = reqwest::Client::builder() .timeout(Duration::from_secs(10)) .build() .context("building an HTTP client")?; let response = client .post(format!("{}/api/login", self.base_url)) // The server requires this on every state-changing request. The // value is irrelevant; a cross-origin browser cannot set it. .header("X-OT-CSRF", "1") .json(&serde_json::json!({ "username": self.username, "password": self.password, "purpose": "device", "device_name": "simulated", "platform": "sim", })) .send() .await .context("posting to /api/login")?; let status = response.status(); let body: serde_json::Value = response .json() .await .context("reading the login response")?; if !status.is_success() { bail!("login failed with {status}: {body}"); } let device = body .get("device") .ok_or_else(|| anyhow::anyhow!("login response has no device credentials"))?; let token_id = device .get("token_id") .and_then(serde_json::Value::as_u64) .ok_or_else(|| anyhow::anyhow!("login response has no token_id"))?; use base64::Engine as _; let key_b64 = device .get("token_key") .and_then(serde_json::Value::as_str) .ok_or_else(|| anyhow::anyhow!("login response has no token_key"))?; let key: Key = base64::engine::general_purpose::STANDARD .decode(key_b64) .context("token_key is not base64")? .try_into() .map_err(|_| anyhow::anyhow!("token_key is not 32 bytes"))?; Ok(Credentials { token_id, k_up: kdf::derive(&key, Direction::Up), k_down: kdf::derive(&key, Direction::Down), }) } /// Send one message and wait briefly for its ACK. /// /// A missing ACK is logged and otherwise ignored: the simulated device has no /// durable queue, and the point of running it is to watch the server, not to /// re-implement the client's retry engine here. async fn send(&self, socket: &UdpSocket, creds: &Credentials, msg: &Message) -> Result<()> { let mut nonce = [0u8; 12]; OsRng.try_fill_bytes(&mut nonce).context("OS RNG failed")?; let datagram = otproto::seal_message(&creds.k_up, creds.token_id, nonce, msg); socket.send(&datagram).await.context("sending a datagram")?; let mut buf = vec![0u8; otproto::MAX_DATAGRAM]; match tokio::time::timeout(Duration::from_secs(2), socket.recv(&mut buf)).await { Ok(Ok(len)) => match otproto::open_message(&creds.k_down, &buf[..len]) { Ok((_, Message::Ack(ack))) => { if ack.nonces.first() != Some(&nonce) { warn!("ACK did not echo the nonce we sent"); } if ack.flags.contains(AckFlags::CONFIG_PENDING) { info!("server has a config update pending"); } } Ok((_, Message::Nack(nack))) => { warn!(reason = ?nack.reason, retry_after_s = nack.retry_after_s, "NACK"); } Ok((_, other)) => warn!(?other, "unexpected reply"), Err(e) => warn!(error = %e, "could not open the reply"), }, Ok(Err(e)) => warn!(error = %e, "recv failed"), Err(_) => warn!("no reply within 2s"), } Ok(()) } }