simulate.rs
⎇
Raw
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
12use std::time::Duration;
13
14use anyhow::{Context, Result, bail};
15use otproto::msg::Direction;
16use otproto::point::Flags;
17use otproto::{AckFlags, Key, Message, Point, kdf};
18use rand::TryRngCore;
19use rand::rngs::OsRng;
20use tokio::net::UdpSocket;
21use tracing::{info, warn};
22
23/// How often the simulated phone reports. Matches the Balanced profile's walking
24/// cadence closely enough to be representative.
25const REPORT_INTERVAL: Duration = Duration::from_secs(5);
26
27/// Roughly walking pace, in units of 1e-7 degrees per report.
28const STEP_E7: i32 = 1_200;
29
30pub struct SimulatedDevice {
31 pub base_url: String,
32 pub username: String,
33 pub password: String,
34 pub udp_addr: String,
35}
36
37struct Credentials {
38 token_id: u64,
39 k_up: Key,
40 k_down: Key,
41}
42
43impl 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