main.rs
| 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 | |
| 15 | mod api; |
| 16 | mod auth; |
| 17 | mod config; |
| 18 | mod db; |
| 19 | mod ingest; |
| 20 | mod keys; |
| 21 | mod limits; |
| 22 | mod polyline; |
| 23 | mod retention; |
| 24 | mod simulate; |
| 25 | mod tiles; |
| 26 | mod udp; |
| 27 | mod web; |
| 28 | mod writer; |
| 29 | |
| 30 | use std::net::SocketAddr; |
| 31 | use std::path::PathBuf; |
| 32 | use std::sync::Arc; |
| 33 | |
| 34 | use anyhow::{Context, Result, bail}; |
| 35 | use tower_sessions::cookie::SameSite; |
| 36 | use tower_sessions::{Expiry, SessionManagerLayer}; |
| 37 | use tower_sessions_sqlx_store::SqliteStore; |
| 38 | use tracing::{info, warn}; |
| 39 | use tracing_subscriber::EnvFilter; |
| 40 | |
| 41 | use crate::api::AppState; |
| 42 | use crate::config::Config; |
| 43 | use crate::db::Db; |
| 44 | use crate::ingest::Ingest; |
| 45 | use crate::keys::KeyVault; |
| 46 | |
| 47 | /// Rolling session lifetime. |
| 48 | const SESSION_DAYS: i64 = 30; |
| 49 | |
| 50 | struct 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 | |
| 60 | fn 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] |
| 101 | async 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. |
| 259 | fn 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 | |
| 267 | async 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)] |
| 289 | mod 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 |