simulate.rs
| 1 | //! `--simulate-device`: an in-process fake phone. |
| 2 | //! |
| 3 | //! It logs in over the real HTTP API, then walks a synthetic route sending real |
| 4 | //! OTP/1 datagrams over a real UDP socket on loopback. Nothing is stubbed: the |
| 5 | //! same codec, the same AEAD, the same ingest path, the same writer. |
| 6 | //! |
| 7 | //! This is the fastest feedback loop for everything downstream of the protocol — |
| 8 | //! the web UI can be built against a moving marker before a line of Kotlin runs — |
| 9 | //! and because it exercises the genuine encoder, a protocol mistake shows up here |
| 10 | //! rather than on a phone. |
| 11 | |
| 12 | use std::time::Duration; |
| 13 | |
| 14 | use anyhow::{Context, Result, bail}; |
| 15 | use otproto::msg::Direction; |
| 16 | use otproto::point::Flags; |
| 17 | use otproto::{AckFlags, Key, Message, Point, kdf}; |
| 18 | use rand::TryRngCore; |
| 19 | use rand::rngs::OsRng; |
| 20 | use tokio::net::UdpSocket; |
| 21 | use tracing::{info, warn}; |
| 22 | |
| 23 | /// How often the simulated phone reports. Matches the Balanced profile's walking |
| 24 | /// cadence closely enough to be representative. |
| 25 | const REPORT_INTERVAL: Duration = Duration::from_secs(5); |
| 26 | |
| 27 | /// Roughly walking pace, in units of 1e-7 degrees per report. |
| 28 | const STEP_E7: i32 = 1_200; |
| 29 | |
| 30 | pub struct SimulatedDevice { |
| 31 | pub base_url: String, |
| 32 | pub username: String, |
| 33 | pub password: String, |
| 34 | pub udp_addr: String, |
| 35 | } |
| 36 | |
| 37 | struct Credentials { |
| 38 | token_id: u64, |
| 39 | k_up: Key, |
| 40 | k_down: Key, |
| 41 | } |
| 42 | |
| 43 | impl SimulatedDevice { |
| 44 | pub fn spawn(self) -> tokio::task::JoinHandle<()> { |
| 45 | tokio::spawn(async move { |
| 46 | if let Err(e) = self.run().await { |
| 47 | warn!(error = %e, "simulated device stopped"); |
| 48 | } |
| 49 | }) |
| 50 | } |
| 51 | |
| 52 | async fn run(self) -> Result<()> { |
| 53 | let creds = self.login().await?; |
| 54 | let socket = UdpSocket::bind("0.0.0.0:0") |
| 55 | .await |
| 56 | .context("binding a client socket")?; |
| 57 | socket |
| 58 | .connect(&self.udp_addr) |
| 59 | .await |
| 60 | .with_context(|| format!("connecting to {}", self.udp_addr))?; |
| 61 | info!(token_id = creds.token_id, addr = %self.udp_addr, "simulated device connected"); |
| 62 | |
| 63 | // Say hello once, exactly as the app does, so the server has a version |
| 64 | // recorded against the token. |
| 65 | let hello = Message::Hello(otproto::Hello { |
| 66 | app_version_code: 0, |
| 67 | os_api_level: 0, |
| 68 | flags: otproto::HelloFlags::FIRST_LAUNCH, |
| 69 | config_version: 1, |
| 70 | }); |
| 71 | self.send(&socket, &creds, &hello).await?; |
| 72 | |
| 73 | // A lap around a small block near the Brandenburg Gate. |
| 74 | let mut lat = 525_200_080i32; |
| 75 | let mut lon = 134_050_000i32; |
| 76 | let mut leg = 0u32; |
| 77 | let mut ticker = tokio::time::interval(REPORT_INTERVAL); |
| 78 | |
| 79 | loop { |
| 80 | ticker.tick().await; |
| 81 | // Four legs of 20 reports each: north, east, south, west. |
| 82 | let (dlat, dlon) = match (leg / 20) % 4 { |
| 83 | 0 => (STEP_E7, 0), |
| 84 | 1 => (0, STEP_E7), |
| 85 | 2 => (-STEP_E7, 0), |
| 86 | _ => (0, -STEP_E7), |
| 87 | }; |
| 88 | lat += dlat; |
| 89 | lon += dlon; |
| 90 | leg = leg.wrapping_add(1); |
| 91 | |
| 92 | let bearing = match (leg / 20) % 4 { |
| 93 | 0 => 0u16, |
| 94 | 1 => 9_000, |
| 95 | 2 => 18_000, |
| 96 | _ => 27_000, |
| 97 | }; |
| 98 | let point = Point { |
| 99 | acc_dm: Some(60), |
| 100 | alt_m: Some(34), |
| 101 | spd_cms: Some(140), |
| 102 | brg_cdeg: Some(bearing), |
| 103 | bat_pct: Some(80u8.saturating_sub((leg / 60) as u8)), |
| 104 | flags: Flags::NONE, |
| 105 | ..Point::new(crate::db::now() as u32, lat, lon) |
| 106 | }; |
| 107 | self.send(&socket, &creds, &Message::Loc(vec![point])) |
| 108 | .await?; |
| 109 | } |
| 110 | } |
| 111 | |
| 112 | async fn login(&self) -> Result<Credentials> { |
| 113 | let client = reqwest::Client::builder() |
| 114 | .timeout(Duration::from_secs(10)) |
| 115 | .build() |
| 116 | .context("building an HTTP client")?; |
| 117 | let response = client |
| 118 | .post(format!("{}/api/login", self.base_url)) |
| 119 | // The server requires this on every state-changing request. The |
| 120 | // value is irrelevant; a cross-origin browser cannot set it. |
| 121 | .header("X-OT-CSRF", "1") |
| 122 | .json(&serde_json::json!({ |
| 123 | "username": self.username, |
| 124 | "password": self.password, |
| 125 | "purpose": "device", |
| 126 | "device_name": "simulated", |
| 127 | "platform": "sim", |
| 128 | })) |
| 129 | .send() |
| 130 | .await |
| 131 | .context("posting to /api/login")?; |
| 132 | |
| 133 | let status = response.status(); |
| 134 | let body: serde_json::Value = response |
| 135 | .json() |
| 136 | .await |
| 137 | .context("reading the login response")?; |
| 138 | if !status.is_success() { |
| 139 | bail!("login failed with {status}: {body}"); |
| 140 | } |
| 141 | |
| 142 | let device = body |
| 143 | .get("device") |
| 144 | .ok_or_else(|| anyhow::anyhow!("login response has no device credentials"))?; |
| 145 | let token_id = device |
| 146 | .get("token_id") |
| 147 | .and_then(serde_json::Value::as_u64) |
| 148 | .ok_or_else(|| anyhow::anyhow!("login response has no token_id"))?; |
| 149 | |
| 150 | use base64::Engine as _; |
| 151 | let key_b64 = device |
| 152 | .get("token_key") |
| 153 | .and_then(serde_json::Value::as_str) |
| 154 | .ok_or_else(|| anyhow::anyhow!("login response has no token_key"))?; |
| 155 | let key: Key = base64::engine::general_purpose::STANDARD |
| 156 | .decode(key_b64) |
| 157 | .context("token_key is not base64")? |
| 158 | .try_into() |
| 159 | .map_err(|_| anyhow::anyhow!("token_key is not 32 bytes"))?; |
| 160 | |
| 161 | Ok(Credentials { |
| 162 | token_id, |
| 163 | k_up: kdf::derive(&key, Direction::Up), |
| 164 | k_down: kdf::derive(&key, Direction::Down), |
| 165 | }) |
| 166 | } |
| 167 | |
| 168 | /// Send one message and wait briefly for its ACK. |
| 169 | /// |
| 170 | /// A missing ACK is logged and otherwise ignored: the simulated device has no |
| 171 | /// durable queue, and the point of running it is to watch the server, not to |
| 172 | /// re-implement the client's retry engine here. |
| 173 | async fn send(&self, socket: &UdpSocket, creds: &Credentials, msg: &Message) -> Result<()> { |
| 174 | let mut nonce = [0u8; 12]; |
| 175 | OsRng.try_fill_bytes(&mut nonce).context("OS RNG failed")?; |
| 176 | let datagram = otproto::seal_message(&creds.k_up, creds.token_id, nonce, msg); |
| 177 | socket.send(&datagram).await.context("sending a datagram")?; |
| 178 | |
| 179 | let mut buf = vec![0u8; otproto::MAX_DATAGRAM]; |
| 180 | match tokio::time::timeout(Duration::from_secs(2), socket.recv(&mut buf)).await { |
| 181 | Ok(Ok(len)) => match otproto::open_message(&creds.k_down, &buf[..len]) { |
| 182 | Ok((_, Message::Ack(ack))) => { |
| 183 | if ack.nonces.first() != Some(&nonce) { |
| 184 | warn!("ACK did not echo the nonce we sent"); |
| 185 | } |
| 186 | if ack.flags.contains(AckFlags::CONFIG_PENDING) { |
| 187 | info!("server has a config update pending"); |
| 188 | } |
| 189 | } |
| 190 | Ok((_, Message::Nack(nack))) => { |
| 191 | warn!(reason = ?nack.reason, retry_after_s = nack.retry_after_s, "NACK"); |
| 192 | } |
| 193 | Ok((_, other)) => warn!(?other, "unexpected reply"), |
| 194 | Err(e) => warn!(error = %e, "could not open the reply"), |
| 195 | }, |
| 196 | Ok(Err(e)) => warn!(error = %e, "recv failed"), |
| 197 | Err(_) => warn!("no reply within 2s"), |
| 198 | } |
| 199 | Ok(()) |
| 200 | } |
| 201 | } |
| 202 |