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