main.rs
⎇
Raw
1//! opentracker server: one binary, one SQLite file.
2//!
3//! Wiring, in dependency order:
4//!
5//! ```text
6//! Db ──▶ writer task ──▶ Ingest ──▶ UDP workers
7//! │ │
8//! └──▶ axum HTTP ─────────┘ (login mints tokens straight into Ingest)
9//! ```
10//!
11//! The one operational trap worth repeating: **HTTP reverse proxies do not forward
12//! UDP.** The HTTP listener belongs behind nginx or Caddy; the UDP port needs its
13//! own firewall rule.
14
15mod api;
16mod auth;
17mod config;
18mod db;
19mod ingest;
20mod keys;
21mod limits;
22mod polyline;
23mod retention;
24mod simulate;
25mod tiles;
26mod udp;
27mod web;
28mod writer;
29
30use std::net::SocketAddr;
31use std::path::PathBuf;
32use std::sync::Arc;
33
34use anyhow::{Context, Result, bail};
35use tower_sessions::cookie::SameSite;
36use tower_sessions::{Expiry, SessionManagerLayer};
37use tower_sessions_sqlx_store::SqliteStore;
38use tracing::{info, warn};
39use tracing_subscriber::EnvFilter;
40
41use crate::api::AppState;
42use crate::config::Config;
43use crate::db::Db;
44use crate::ingest::Ingest;
45use crate::keys::KeyVault;
46
47/// Rolling session lifetime.
48const SESSION_DAYS: i64 = 30;
49
50struct Args {
51 config: Option<PathBuf>,
52 /// Relaxes the cookie's `Secure` requirement so the Vite dev server works
53 /// over plain HTTP on localhost. Never for production.
54 dev: bool,
55 simulate_device: bool,
56 /// `--create-admin <user> <password>`: bootstrap the first account.
57 create_admin: Option<(String, String)>,
58}
59
60fn parse_args() -> Result<Args> {
61 let mut args = Args {
62 config: None,
63 dev: false,
64 simulate_device: false,
65 create_admin: None,
66 };
67 let mut it = std::env::args().skip(1);
68 while let Some(arg) = it.next() {
69 match arg.as_str() {
70 "--config" | "-c" => {
71 args.config = Some(PathBuf::from(it.next().context("--config needs a path")?));
72 }
73 "--dev" => args.dev = true,
74 "--simulate-device" => args.simulate_device = true,
75 "--create-admin" => {
76 let user = it.next().context("--create-admin needs a username")?;
77 let pass = it.next().context("--create-admin needs a password")?;
78 args.create_admin = Some((user, pass));
79 }
80 "--help" | "-h" => {
81 println!(
82 "opentracker {}\n\n\
83 USAGE:\n \
84 otserver [--config FILE] [--dev] [--simulate-device]\n \
85 otserver --create-admin USERNAME PASSWORD\n\n\
86 ENVIRONMENT:\n \
87 OT_SECRET_KEY required; 32 random bytes, base64, or a path to them\n \
88 OT_SECRET_KEY_OLD accepted for unwrapping during a key rotation\n \
89 OT_* override any config field (see config.rs)\n",
90 env!("CARGO_PKG_VERSION")
91 );
92 std::process::exit(0);
93 }
94 other => bail!("unknown argument {other:?} (try --help)"),
95 }
96 }
97 Ok(args)
98}
99
100#[tokio::main]
101async fn main() -> Result<()> {
102 tracing_subscriber::fmt()
103 .with_env_filter(
104 EnvFilter::try_from_env("OT_LOG")
105 .unwrap_or_else(|_| EnvFilter::new("info,otserver=debug")),
106 )
107 .init();
108
109 let args = parse_args()?;
110 let cfg = Config::load(args.config.as_deref())?;
111 cfg.validate()?;
112
113 let vault = KeyVault::from_env()?;
114 let db = Db::open(&cfg.db_path).await?;
115 info!(path = %cfg.db_path.display(), "database ready");
116
117 let (writer, writer_task) = writer::spawn(db.write.clone());
118 // Handing over the master is what enables replies to datagrams naming tokens
119 // this server has no record of. Withheld unless config asks for it, so the
120 // one reflective path cannot be switched on by accident.
121 let revocation_master = cfg.revocation_notices.then(|| vault.revocation_master());
122 if revocation_master.is_some() {
123 info!(
124 "revocation_notices enabled: an unknown token draws a rate-limited 38-byte reply. \
125 See config.rs for the trade."
126 );
127 }
128 let ingest = Arc::new(Ingest::new(
129 writer.clone(),
130 cfg.ts_window_days * 86_400,
131 revocation_master,
132 ));
133
134 // One-shot bootstrap, before anything starts listening.
135 if let Some((username, password)) = args.create_admin {
136 let hash = auth::hash_password(&cfg, password).await?;
137 let at = db::now();
138 sqlx::query(
139 "INSERT INTO users (username, pw_hash, display_name, is_admin, created_at, pw_changed_at) \
140 VALUES (?, ?, ?, 1, ?, ?)",
141 )
142 .bind(&username)
143 .bind(hash)
144 .bind(&username)
145 .bind(at)
146 .bind(at)
147 .execute(&db.write)
148 .await
149 .with_context(|| format!("creating admin {username}"))?;
150 println!("created admin account {username}");
151 db.checkpoint().await?;
152 db.close().await;
153 return Ok(());
154 }
155
156 let loaded = auth::load_tokens(&db.read, &vault, &ingest).await?;
157 info!(
158 tokens = loaded,
159 "loaded device tokens into the ingest cache"
160 );
161
162 // Sessions live in the same SQLite file, on the writer pool. They are
163 // low-traffic enough not to disturb the batching that exists to protect the
164 // position write path.
165 let session_store = SqliteStore::new(db.write.clone());
166 session_store
167 .migrate()
168 .await
169 .context("migrating the session store")?;
170
171 if args.dev {
172 warn!("--dev: session cookies will not require HTTPS. Never use this in production.");
173 }
174 let session_layer = SessionManagerLayer::new(session_store)
175 .with_name(if args.dev { "otsid" } else { "__Host-otsid" })
176 .with_http_only(true)
177 .with_secure(!args.dev)
178 .with_same_site(SameSite::Lax)
179 .with_expiry(Expiry::OnInactivity(time::Duration::days(SESSION_DAYS)));
180
181 let http_addr = cfg.http_addr;
182 let udp_addr = cfg.udp_addr;
183 let workers = cfg.udp_worker_count();
184 let retention_days = cfg.retention_days;
185 let token_stale_days = cfg.token_stale_days;
186
187 let throttle = Arc::new(auth::LoginThrottle::default());
188 let state: api::Shared = Arc::new(AppState {
189 db,
190 cfg,
191 vault,
192 ingest: Arc::clone(&ingest),
193 writer,
194 throttle: Arc::clone(&throttle),
195 });
196
197 let udp_tasks = udp::spawn(udp_addr, workers, Arc::clone(&ingest))?;
198
199 let retention_task = retention::Retention {
200 pool: state.db.write.clone(),
201 retention_days,
202 token_stale_days,
203 ingest: Arc::clone(&ingest),
204 }
205 .spawn();
206 let limits_task = retention::spawn_limits_gc(Arc::clone(&ingest), throttle);
207 let eviction_task = tiles::spawn_eviction(
208 state.db.write.clone(),
209 state.cfg.cache_dir.clone(),
210 state.cfg.max_cache_bytes,
211 );
212
213 let app = api::router(Arc::clone(&state)).layer(session_layer);
214
215 let listener = tokio::net::TcpListener::bind(http_addr)
216 .await
217 .with_context(|| format!("binding {http_addr}"))?;
218 info!(addr = %http_addr, "HTTP listener started (put a TLS-terminating proxy in front)");
219
220 if args.simulate_device {
221 // After the listeners are up, so the login cannot race them.
222 simulate::SimulatedDevice {
223 base_url: format!("http://{http_addr}"),
224 username: std::env::var("OT_SIM_USER").unwrap_or_else(|_| "sim".into()),
225 password: std::env::var("OT_SIM_PASSWORD").unwrap_or_else(|_| "simsimsimsim".into()),
226 udp_addr: loopback_of(udp_addr),
227 }
228 .spawn();
229 }
230
231 axum::serve(
232 listener,
233 app.into_make_service_with_connect_info::<SocketAddr>(),
234 )
235 .with_graceful_shutdown(shutdown_signal())
236 .await
237 .context("HTTP server failed")?;
238
239 // Stop the periodic tasks and the receive loops, then let the writer drain so
240 // nothing already acknowledged is lost, then checkpoint so the .db file on
241 // disk is self-contained.
242 info!("shutting down");
243 retention_task.abort();
244 limits_task.abort();
245 eviction_task.abort();
246 for task in udp_tasks {
247 task.abort();
248 }
249 if let Err(e) = state.db.checkpoint().await {
250 warn!(error = %e, "WAL checkpoint failed");
251 }
252 drop(state);
253 let _ = tokio::time::timeout(std::time::Duration::from_secs(5), writer_task).await;
254 Ok(())
255}
256
257/// The simulated device must talk to a routable address; the listener is usually
258/// bound to the wildcard, which cannot be a destination.
259fn loopback_of(addr: SocketAddr) -> String {
260 if addr.ip().is_unspecified() {
261 format!("127.0.0.1:{}", addr.port())
262 } else {
263 addr.to_string()
264 }
265}
266
267async fn shutdown_signal() {
268 let ctrl_c = async {
269 tokio::signal::ctrl_c().await.ok();
270 };
271 #[cfg(unix)]
272 let terminate = async {
273 if let Ok(mut sig) =
274 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
275 {
276 sig.recv().await;
277 }
278 };
279 #[cfg(not(unix))]
280 let terminate = std::future::pending::<()>();
281
282 tokio::select! {
283 () = ctrl_c => {},
284 () = terminate => {},
285 }
286}
287
288#[cfg(test)]
289mod tests {
290 use super::*;
291
292 #[test]
293 fn the_wildcard_address_is_rewritten_for_the_simulated_device() {
294 assert_eq!(
295 loopback_of("0.0.0.0:7373".parse().expect("literal")),
296 "127.0.0.1:7373"
297 );
298 assert_eq!(
299 loopback_of("192.0.2.1:7373".parse().expect("literal")),
300 "192.0.2.1:7373"
301 );
302 }
303}
304